亚洲乱码中文字幕综合,中国熟女仑乱hd,亚洲精品乱拍国产一区二区三区,一本大道卡一卡二卡三乱码全集资源,又粗又黄又硬又爽的免费视频

SpringBoot集成Kafka的步驟

 更新時間:2021年01月06日 11:43:01   作者:Aska小強(qiáng)  
這篇文章主要介紹了SpringBoot集成Kafka的步驟,幫助大家更好的理解和使用SpringBoot,感興趣的朋友可以了解下

SpringBoot集成Kafka

本篇主要講解SpringBoot 如何集成Kafka ,并且簡單的 編寫了一個Demo 來測試 發(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 它會引入 spring-kafka 2.5.6 RELEASE ,對應(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)境 并且啟動 它們

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的地址,多個地址用逗號分隔
  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的值,客戶端會將發(fā)送失敗的記錄重新發(fā)送
   batch-size: 16384 #當(dāng)將多個記錄被發(fā)送到同一個分區(qū)時, Producer 將嘗試將記錄組合到更少的請求中。這有助于提升客戶端和服務(wù)器端的性能。這個配置控制一個批次的默認(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));
    //對 topic = hello2 的發(fā)送消息
    kafkaTemplate.send("hello2",gson.toJson(message));
  }

}

2.5 定義KafkaConsumer

在監(jiān)聽的方法上通過注解配置一個監(jiān)聽器即可,另外就是指定需要監(jiān)聽的topic
kafka的消息再接收端會被封裝成ConsumerRecord對象返回,它內(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.測試 效果

提供一個 Http接口調(diào)用 KafkaSender 去發(fā)送消息

3.1 提供Http 測試接口

@RestController
@Slf4j
public class TestController {


  @Autowired
  private KafkaSender kafkaSender;


  @GetMapping("sendMessage/{msg}")
  public void sendMessage(@PathVariable("msg") String msg){
    kafkaSender.send(msg);
  }
}

3.2 啟動項(xiàng)目

監(jiān)聽8080 端口

KafkaMessageListenerContainer中有 consumer group = myGroup 有一個 監(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的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Java程序員必須知道的5個JVM命令行標(biāo)志

    Java程序員必須知道的5個JVM命令行標(biāo)志

    這篇文章主要介紹了每個Java程序員必須知道的5個JVM命令行標(biāo)志,需要的朋友可以參考下
    2015-03-03
  • java實(shí)現(xiàn)文件斷點(diǎn)續(xù)傳下載功能

    java實(shí)現(xiàn)文件斷點(diǎn)續(xù)傳下載功能

    這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)文件斷點(diǎn)續(xù)傳下載功能的具體代碼,感興趣的小伙伴們可以參考一下
    2016-05-05
  • @DynamicUpdate //自動更新updatetime的問題

    @DynamicUpdate //自動更新updatetime的問題

    這篇文章主要介紹了@DynamicUpdate //自動更新updatetime的問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-07-07
  • Spring中自動裝配的4種方式

    Spring中自動裝配的4種方式

    今天小編就為大家分享一篇關(guān)于Spring中自動裝配的4種方式,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2019-01-01
  • Java擴(kuò)展庫RxJava的基本結(jié)構(gòu)與適用場景小結(jié)

    Java擴(kuò)展庫RxJava的基本結(jié)構(gòu)與適用場景小結(jié)

    RxJava(GitHub: https://github.com/ReactiveX/RxJava)能夠幫助Java進(jìn)行異步與事務(wù)驅(qū)動的程序編寫,這里我們來作一個Java擴(kuò)展庫RxJava的基本結(jié)構(gòu)與適用場景小結(jié),剛接觸RxJava的同學(xué)不妨看一下^^
    2016-06-06
  • 30分鐘入門Java8之lambda表達(dá)式學(xué)習(xí)

    30分鐘入門Java8之lambda表達(dá)式學(xué)習(xí)

    本篇文章主要介紹了30分鐘入門Java8之lambda表達(dá)式學(xué)習(xí),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-04-04
  • Java中Lambda表達(dá)式基礎(chǔ)及使用

    Java中Lambda表達(dá)式基礎(chǔ)及使用

    這篇文章主要介紹了Lambda 是JDK 8 的重要新特性。它允許把函數(shù)作為一個方法的參數(shù)(函數(shù)作為參數(shù)傳遞進(jìn)方法中),使用 Lambda 表達(dá)式可以使代碼變的更加簡潔緊湊,使Java代碼更加優(yōu)雅,感興趣的小伙伴一起來學(xué)習(xí)吧
    2021-08-08
  • Java之next()、nextLine()區(qū)別及問題解決

    Java之next()、nextLine()區(qū)別及問題解決

    這篇文章主要介紹了Java之next()、nextLine()區(qū)別及問題解決,本篇文章通過簡要的案例,講解了該項(xiàng)技術(shù)的了解與使用,以下就是詳細(xì)內(nèi)容,需要的朋友可以參考下
    2021-08-08
  • Java日常練習(xí)題,每天進(jìn)步一點(diǎn)點(diǎn)(43)

    Java日常練習(xí)題,每天進(jìn)步一點(diǎn)點(diǎn)(43)

    下面小編就為大家?guī)硪黄狫ava基礎(chǔ)的幾道練習(xí)題(分享)。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧,希望可以幫到你
    2021-07-07
  • 用IDEA如何打開文件夾

    用IDEA如何打開文件夾

    這篇文章主要介紹了用IDEA如何打開文件夾問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2023-09-09

最新評論