最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

spring-kafka使消費(fèi)者動(dòng)態(tài)訂閱新增的topic問(wèn)題

 更新時(shí)間:2022年12月27日 11:31:56   作者:DayDayUp丶  
這篇文章主要介紹了spring-kafka使消費(fèi)者動(dòng)態(tài)訂閱新增的topic問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教

一、前言

在Java中使用kafka,方式很多,例如:

  • 直接使用kafka-clients這類原生的API;
  • 也可以使用Spring對(duì)其的包裝API,即spring-kafka,同其它包裝API一樣(如JdbcTemplate、RestTemplate、RedisTemplate等等),KafkaTemplate是其生產(chǎn)者核心類,KafkaListener是其消費(fèi)者核心注解;
  • 也有包裝地更加抽象的SpringCloudStream等。

這里討論的話題是,如何在spring-kafka中,使得一個(gè)消費(fèi)者可以動(dòng)態(tài)訂閱新增的topic?

本文不討論利用SpringCloudConfig或Apollo等分布式配置中心,利用@RefreshScope的方式來(lái)達(dá)到目的,這種方式有點(diǎn)殺雞用牛刀,也會(huì)增加系統(tǒng)復(fù)雜度和維護(hù)成本。

我的環(huán)境:jdk 1.8,Spring 2.1.3.RELEASE,kafka_2.12-2.3.0單節(jié)點(diǎn)。

二、需求分析

上面已經(jīng)提到,spring-kafka通過(guò) @KafkaListener 的方式配置訂閱的topic,最常用的屬性可能是 topics,而要實(shí)現(xiàn)本文的需求,就要使用另一個(gè)屬性 topicPattern,查看它的屬性說(shuō)明:

The topic pattern for this listener. 
The entries can be 'topic pattern', a'property-placeholder key' or an 'expression'. 
The framework will create acontainer that subscribes to all topics matching the specified pattern to getdynamically assigned partitions. 
The pattern matching will be performedperiodically against topics existing at the time of check. 
An expression mustbe resolved to the topic pattern (String or Pattern result types are supported). 

將其翻譯過(guò)來(lái):

此偵聽器的主題模式。條目可以是“主題模式”,“屬性占位符鍵”或“表達(dá)式”。
該框架將創(chuàng)建一個(gè)容器,該容器訂閱與指定模式匹配的所有主題以獲取動(dòng)態(tài)分配的分區(qū)。
模式匹配將針對(duì)檢查時(shí)存在的主題【定期執(zhí)行】。
表達(dá)式必須解析為主題模式(支持字符串或模式結(jié)果類型)。

注意:從說(shuō)明信息來(lái)看,topicPattern 已經(jīng)可以做到定期檢查topic列表,然后將新加入的topic分配至某個(gè)消費(fèi)者。

下面列出消費(fèi)端的核心測(cè)試代碼:

@Component
public class SinkConsumer {
    @KafkaListener(topicPattern = "test_topic2.*")
    public void listen2(ConsumerRecord<?, ?> record) throws Exception {
        System.out.printf("topic2.* = %s, offset = %d, value = %s \n", record.topic(), record.offset(), record.value());
    }
}

代碼實(shí)現(xiàn)很簡(jiǎn)潔,就是期待我們新增一個(gè)符合 topicPattern 的topic后,spring-kafka能否自動(dòng)為新建的topic分配到此目標(biāo)消費(fèi)者。

三、測(cè)試運(yùn)行

3.1 啟動(dòng)消費(fèi)者服務(wù)

配置文件中,spring該配的配,kafka該配的配,接著啟動(dòng)即可。

3.2 新建topic

新建 test_topic2_3,剛創(chuàng)建完不能立刻分配到目標(biāo)消費(fèi)者,從 topicPattern 的注釋得知spring-kafka會(huì)定期掃描topic列表,我們要給它幾分鐘等待掃描到新topic,并為它成功分配到目標(biāo)消費(fèi)者后,再去發(fā)送第一條消息(所以可以先去洗個(gè)手,此時(shí)19:02)。

3.3 等待topic被分配到消費(fèi)者

洗手期間的控制臺(tái)日志提示:已為新建的 test_topic2_3 分配到我們的目標(biāo)消費(fèi)者,并將offset設(shè)置到起始位置0,日志如下:

2019-11-15 19:05:12.958  INFO 7768 --- [ntainer#1-0-C-1] o.a.k.c.c.internals.ConsumerCoordinator  : [Consumer clientId=consumer-3, groupId=test] Revoking previously assigned partitions [test_topic2_2-0, test_topic2_1-0]
2019-11-15 19:05:12.958  INFO 7768 --- [ntainer#1-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : partitions revoked: [test_topic2_2-0, test_topic2_1-0]
2019-11-15 19:05:12.958  INFO 7768 --- [ntainer#1-0-C-1] o.a.k.c.c.internals.AbstractCoordinator  : [Consumer clientId=consumer-3, groupId=test] (Re-)joining group
2019-11-15 19:05:15.757  INFO 7768 --- [ntainer#0-0-C-1] o.a.k.c.c.internals.AbstractCoordinator  : [Consumer clientId=consumer-2, groupId=test] Attempt to heartbeat failed since group is rebalancing
2019-11-15 19:05:15.761  INFO 7768 --- [ntainer#0-0-C-1] o.a.k.c.c.internals.ConsumerCoordinator  : [Consumer clientId=consumer-2, groupId=test] Revoking previously assigned partitions [test_topic-0]
2019-11-15 19:05:15.762  INFO 7768 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : partitions revoked: [test_topic-0]
2019-11-15 19:05:15.762  INFO 7768 --- [ntainer#0-0-C-1] o.a.k.c.c.internals.AbstractCoordinator  : [Consumer clientId=consumer-2, groupId=test] (Re-)joining group
2019-11-15 19:05:16.025  INFO 7768 --- [ntainer#1-0-C-1] o.a.k.c.c.internals.AbstractCoordinator  : [Consumer clientId=consumer-3, groupId=test] Successfully joined group with generation 6
2019-11-15 19:05:16.025  INFO 7768 --- [ntainer#0-0-C-1] o.a.k.c.c.internals.AbstractCoordinator  : [Consumer clientId=consumer-2, groupId=test] Successfully joined group with generation 6
2019-11-15 19:05:16.026  INFO 7768 --- [ntainer#1-0-C-1] o.a.k.c.c.internals.ConsumerCoordinator  : [Consumer clientId=consumer-3, groupId=test] Setting newly assigned partitions [test_topic2_2-0, test_topic2_3-0, test_topic2_1-0]
2019-11-15 19:05:16.026  INFO 7768 --- [ntainer#0-0-C-1] o.a.k.c.c.internals.ConsumerCoordinator  : [Consumer clientId=consumer-2, groupId=test] Setting newly assigned partitions [test_topic-0]
2019-11-15 19:05:16.028  INFO 7768 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : partitions assigned: [test_topic-0]
2019-11-15 19:05:16.032  INFO 7768 --- [ntainer#1-0-C-1] o.a.k.c.consumer.internals.Fetcher       : [Consumer clientId=consumer-3, groupId=test] Resetting offset for partition test_topic2_3-0 to offset 0.
2019-11-15 19:05:16.032  INFO 7768 --- [ntainer#1-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : partitions assigned: [test_topic2_2-0, test_topic2_3-0, test_topic2_1-0]

3.4 發(fā)送第一條消息

洗手完畢,看到3.3小節(jié)里的日志,然后確認(rèn)成功分配到目標(biāo)消費(fèi)者,且offset被設(shè)為0之后,發(fā)送第一條消息【我是第1個(gè)test_topic2_3的消息】,控制臺(tái)日志打印出此消息信息,代表成功消費(fèi):

topic2.* = test_topic2_3, offset = 0, value = {"date":"2019-11-15 19:11:13","msg":"我是第1個(gè)test_topic2_3的消息"} 

3.5 注意事項(xiàng)

若不等到offset被設(shè)為0之后,過(guò)早發(fā)送消息,則會(huì)在消費(fèi)端丟失過(guò)早發(fā)送的消息,并且當(dāng)spring-kafka自動(dòng)設(shè)置offset的時(shí)候,日志提示,offset被設(shè)置為1,而不是起始位置0:

INFO o.a.k.c.consumer.internals.Fetcher       : [Consumer clientId=consumer-3, groupId=test] Resetting offset for partition test_topic2_1-0 to offset 1.

在上面的3.1至3.4的整個(gè)過(guò)程中,可能會(huì)日志警告,代表暫時(shí)不能為新增的topic分配到目標(biāo)消費(fèi)者:

WARN o.a.k.c.c.internals.ConsumerCoordinator  : [Consumer clientId=consumer-2, groupId=test] The following subscribed topics are not assigned to any members: [test_topic2_3] 

所以只需等待日志提示可以成功分配到目標(biāo)消費(fèi)者,且offset被設(shè)為0之后,即可發(fā)送第一條消息。

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • mybatis多個(gè)區(qū)間處理方式(雙foreach循環(huán))

    mybatis多個(gè)區(qū)間處理方式(雙foreach循環(huán))

    這篇文章主要介紹了mybatis多個(gè)區(qū)間處理方式(雙foreach循環(huán)),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • 關(guān)于Java反射給泛型集合賦值問(wèn)題

    關(guān)于Java反射給泛型集合賦值問(wèn)題

    這篇文章主要介紹了Java反射給泛型集合賦值,需要的朋友可以參考下
    2022-01-01
  • Java開啟線程的四種方法案例詳解

    Java開啟線程的四種方法案例詳解

    這篇文章主要介紹了Java開啟線程的四種方法,本文結(jié)合實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2023-02-02
  • SpringBoot中使用tkMapper的方法詳解

    SpringBoot中使用tkMapper的方法詳解

    這篇文章主要介紹了SpringBoot中使用tkMapper的方法詳解
    2022-11-11
  • java簡(jiǎn)單實(shí)現(xiàn)多線程及線程池實(shí)例詳解

    java簡(jiǎn)單實(shí)現(xiàn)多線程及線程池實(shí)例詳解

    這篇文章主要為大家詳細(xì)介紹了java簡(jiǎn)單實(shí)現(xiàn)多線程,及java爬蟲使用線程池實(shí)例,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-03-03
  • Java Scanner對(duì)象中hasNext()與next()方法的使用

    Java Scanner對(duì)象中hasNext()與next()方法的使用

    這篇文章主要介紹了Java Scanner對(duì)象中hasNext()與next()方法的使用,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-10-10
  • Java反射中java.beans包學(xué)習(xí)總結(jié)

    Java反射中java.beans包學(xué)習(xí)總結(jié)

    本篇文章通過(guò)學(xué)習(xí)Java反射中java.beans包,吧知識(shí)點(diǎn)做了總結(jié),并把相關(guān)內(nèi)容做了關(guān)聯(lián),對(duì)此有需要的朋友可以學(xué)習(xí)參考下。
    2018-02-02
  • 詳解Java多線程和IO流的應(yīng)用

    詳解Java多線程和IO流的應(yīng)用

    這篇文章主要介紹了詳解Java多線程和IO流的應(yīng)用,無(wú)論是本地文件復(fù)制,還是網(wǎng)絡(luò)多線程下載,對(duì)于流的使用都是一樣的,需要的朋友可以參考下
    2023-04-04
  • 使用maven自定義插件開發(fā)

    使用maven自定義插件開發(fā)

    這篇文章主要介紹了使用maven自定義插件開發(fā),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • Java8之Lambda表達(dá)式使用解讀

    Java8之Lambda表達(dá)式使用解讀

    這篇文章主要介紹了Java8之Lambda表達(dá)式使用解讀,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-11-11

最新評(píng)論

开封市| 鄂托克前旗| 德兴市| 云和县| 无极县| 淄博市| 平利县| 宁陵县| 准格尔旗| 故城县| 宝清县| 虞城县| 古浪县| 大埔区| 绩溪县| 江阴市| 东山县| 抚顺市| 望都县| 六枝特区| 光山县| 晋城| 桑日县| 双鸭山市| 陇西县| 宁夏| 孝昌县| 临邑县| 灵丘县| 高淳县| 新干县| 黔西县| 资讯 | 遂昌县| 通道| 紫金县| 阿城市| 武鸣县| 丰宁| 沛县| 铁岭市|