SpringBoot集成Kafka的步驟
SpringBoot集成Kafka
本篇主要講解SpringBoot 如何集成Kafka ,并且簡單的 編寫了一個(gè)Demo 來測(cè)試 發(fā)送和消費(fèi)功能
前言
選擇的版本如下:
springboot : 2.3.4.RELEASE
spring-kafka : 2.5.6.RELEASE
kafka : 2.5.1
zookeeper : 3.4.14
本Demo 使用的是 SpringBoot 比較高的版本 SpringBoot 2.3.4.RELEASE 它會(huì)引入 spring-kafka 2.5.6 RELEASE ,對(duì)應(yīng)了版本關(guān)系中的
Spring Boot 2.3 users should use 2.5.x (Boot dependency management will use the correct version).
spring和 kafka 的版本 關(guān)系
https://spring.io/projects/sp...
1.搭建Kafka 和 Zookeeper 環(huán)境
搭建kafka 和 zookeeper 環(huán)境 并且啟動(dòng) 它們
2.創(chuàng)建Demo 項(xiàng)目引入spring-kafka
2.1 pom 文件
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>com.google.code.gson</groupId> <artifactId>gson</artifactId> </dependency>
2.2 配置application.yml
spring: kafka: bootstrap-servers: 192.168.25.6:9092 #bootstrap-servers:連接kafka的地址,多個(gè)地址用逗號(hào)分隔 consumer: group-id: myGroup enable-auto-commit: true auto-commit-interval: 100ms properties: session.timeout.ms: 15000 key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest producer: retries: 0 #若設(shè)置大于0的值,客戶端會(huì)將發(fā)送失敗的記錄重新發(fā)送 batch-size: 16384 #當(dāng)將多個(gè)記錄被發(fā)送到同一個(gè)分區(qū)時(shí), Producer 將嘗試將記錄組合到更少的請(qǐng)求中。這有助于提升客戶端和服務(wù)器端的性能。這個(gè)配置控制一個(gè)批次的默認(rèn)大?。ㄒ宰止?jié)為單位)。16384是缺省的配置 buffer-memory: 33554432 #Producer 用來緩沖等待被發(fā)送到服務(wù)器的記錄的總字節(jié)數(shù),33554432是缺省配置 key-serializer: org.apache.kafka.common.serialization.StringSerializer #關(guān)鍵字的序列化類 value-serializer: org.apache.kafka.common.serialization.StringSerializer #值的序列化類
2.3 定義消息體Message
/** * @author johnny * @create 2020-09-23 上午9:21 **/ @Data public class Message { private Long id; private String msg; private Date sendTime; }
2.4 定義KafkaSender
主要利用 KafkaTemplate 來發(fā)送消息 ,將消息封裝成Message 并且進(jìn)行 轉(zhuǎn)化成Json串 發(fā)送到Kafka中
@Component @Slf4j public class KafkaSender { private final KafkaTemplate<String, String> kafkaTemplate; //構(gòu)造器方式注入 kafkaTemplate public KafkaSender(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } private Gson gson = new GsonBuilder().create(); public void send(String msg) { Message message = new Message(); message.setId(System.currentTimeMillis()); message.setMsg(msg); message.setSendTime(new Date()); log.info("【++++++++++++++++++ message :{}】", gson.toJson(message)); //對(duì) topic = hello2 的發(fā)送消息 kafkaTemplate.send("hello2",gson.toJson(message)); } }
2.5 定義KafkaConsumer
在監(jiān)聽的方法上通過注解配置一個(gè)監(jiān)聽器即可,另外就是指定需要監(jiān)聽的topic
kafka的消息再接收端會(huì)被封裝成ConsumerRecord對(duì)象返回,它內(nèi)部的value屬性就是實(shí)際的消息。
@Component @Slf4j public class KafkaConsumer { @KafkaListener(topics = {"hello2"}) public void listen(ConsumerRecord<?, ?> record) { Optional.ofNullable(record.value()) .ifPresent(message -> { log.info("【+++++++++++++++++ record = {} 】", record); log.info("【+++++++++++++++++ message = {}】", message); }); } }
3.測(cè)試 效果
提供一個(gè) Http接口調(diào)用 KafkaSender 去發(fā)送消息
3.1 提供Http 測(cè)試接口
@RestController @Slf4j public class TestController { @Autowired private KafkaSender kafkaSender; @GetMapping("sendMessage/{msg}") public void sendMessage(@PathVariable("msg") String msg){ kafkaSender.send(msg); } }
3.2 啟動(dòng)項(xiàng)目
監(jiān)聽8080 端口
KafkaMessageListenerContainer中有 consumer group = myGroup 有一個(gè) 監(jiān)聽 hello2-0 topic 的 消費(fèi)者
3.3 調(diào)用Http接口
http://localhost:8080/sendMessage/KafkaTestMsg
至此 SpringBoot集成Kafka 結(jié)束 。。
以上就是SpringBoot集成Kafka的步驟的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot集成Kafka的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
- 如何使用SpringBoot集成Kafka實(shí)現(xiàn)用戶數(shù)據(jù)變更后發(fā)送消息
- SpringBoot3集成Kafka的方法詳解
- SpringBoot集成Kafka 配置工具類的詳細(xì)代碼
- springboot集成kafka消費(fèi)手動(dòng)啟動(dòng)停止操作
- Springboot集成kafka高級(jí)應(yīng)用實(shí)戰(zhàn)分享
- Springboot 2.x集成kafka 2.2.0的示例代碼
- SpringBoot集成kafka全面實(shí)戰(zhàn)記錄
- Springboot集成Kafka實(shí)現(xiàn)producer和consumer的示例代碼
- SpringBoot集成Kafka的實(shí)現(xiàn)示例
相關(guān)文章
Java程序員必須知道的5個(gè)JVM命令行標(biāo)志
這篇文章主要介紹了每個(gè)Java程序員必須知道的5個(gè)JVM命令行標(biāo)志,需要的朋友可以參考下2015-03-03java實(shí)現(xiàn)文件斷點(diǎn)續(xù)傳下載功能
這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)文件斷點(diǎn)續(xù)傳下載功能的具體代碼,感興趣的小伙伴們可以參考一下2016-05-05@DynamicUpdate //自動(dòng)更新updatetime的問題
這篇文章主要介紹了@DynamicUpdate //自動(dòng)更新updatetime的問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-07-07Java擴(kuò)展庫RxJava的基本結(jié)構(gòu)與適用場(chǎng)景小結(jié)
RxJava(GitHub: https://github.com/ReactiveX/RxJava)能夠幫助Java進(jìn)行異步與事務(wù)驅(qū)動(dòng)的程序編寫,這里我們來作一個(gè)Java擴(kuò)展庫RxJava的基本結(jié)構(gòu)與適用場(chǎng)景小結(jié),剛接觸RxJava的同學(xué)不妨看一下^^2016-06-0630分鐘入門Java8之lambda表達(dá)式學(xué)習(xí)
本篇文章主要介紹了30分鐘入門Java8之lambda表達(dá)式學(xué)習(xí),小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧2017-04-04Java之next()、nextLine()區(qū)別及問題解決
這篇文章主要介紹了Java之next()、nextLine()區(qū)別及問題解決,本篇文章通過簡要的案例,講解了該項(xiàng)技術(shù)的了解與使用,以下就是詳細(xì)內(nèi)容,需要的朋友可以參考下2021-08-08Java日常練習(xí)題,每天進(jìn)步一點(diǎn)點(diǎn)(43)
下面小編就為大家?guī)硪黄狫ava基礎(chǔ)的幾道練習(xí)題(分享)。小編覺得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧,希望可以幫到你2021-07-07