亚洲激情专区-91九色丨porny丨老师-久久久久久久女国产乱让韩-国产精品午夜小视频观看

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

spring kafka怎么實現消費者動態訂閱新增的topic

發布時間:2022-12-28 09:30:50 來源:億速云 閱讀:181 作者:iii 欄目:開發技術

這篇文章主要介紹了spring kafka怎么實現消費者動態訂閱新增的topic的相關知識,內容詳細易懂,操作簡單快捷,具有一定借鑒價值,相信大家閱讀完這篇spring kafka怎么實現消費者動態訂閱新增的topic文章都會有所收獲,下面我們一起來看看吧。

    一、前言

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

    • 直接使用kafka-clients這類原生的API;

    • 也可以使用Spring對其的包裝API,即spring-kafka,同其它包裝API一樣(如JdbcTemplate、RestTemplate、RedisTemplate等等),KafkaTemplate是其生產者核心類,KafkaListener是其消費者核心注解;

    • 也有包裝地更加抽象的SpringCloudStream等。

    這里討論的話題是,如何在spring-kafka中,使得一個消費者可以動態訂閱新增的topic?

    本文不討論利用SpringCloudConfig或Apollo等分布式配置中心,利用@RefreshScope的方式來達到目的,這種方式有點殺雞用牛刀,也會增加系統復雜度和維護成本。

    我的環境:jdk 1.8,Spring 2.1.3.RELEASE,kafka_2.12-2.3.0單節點。

    二、需求分析

    上面已經提到,spring-kafka通過 @KafkaListener 的方式配置訂閱的topic,最常用的屬性可能是 topics,而要實現本文的需求,就要使用另一個屬性 topicPattern,查看它的屬性說明:

    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). 

    將其翻譯過來:

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

    注意:從說明信息來看,topicPattern 已經可以做到定期檢查topic列表,然后將新加入的topic分配至某個消費者。

    下面列出消費端的核心測試代碼:

    @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());
        }
    }

    代碼實現很簡潔,就是期待我們新增一個符合 topicPattern 的topic后,spring-kafka能否自動為新建的topic分配到此目標消費者。

    三、測試運行

    3.1 啟動消費者服務

    配置文件中,spring該配的配,kafka該配的配,接著啟動即可。

    3.2 新建topic

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

    3.3 等待topic被分配到消費者

    洗手期間的控制臺日志提示:已為新建的 test_topic2_3 分配到我們的目標消費者,并將offset設置到起始位置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 發送第一條消息

    洗手完畢,看到3.3小節里的日志,然后確認成功分配到目標消費者,且offset被設為0之后,發送第一條消息【我是第1個test_topic2_3的消息】,控制臺日志打印出此消息信息,代表成功消費:

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

    3.5 注意事項

    若不等到offset被設為0之后,過早發送消息,則會在消費端丟失過早發送的消息,并且當spring-kafka自動設置offset的時候,日志提示,offset被設置為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的整個過程中,可能會日志警告,代表暫時不能為新增的topic分配到目標消費者:

    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] 

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

    關于“spring kafka怎么實現消費者動態訂閱新增的topic”這篇文章的內容就介紹到這里,感謝各位的閱讀!相信大家對“spring kafka怎么實現消費者動態訂閱新增的topic”知識都有一定的了解,大家如果還想學習更多知識,歡迎關注億速云行業資訊頻道。

    向AI問一下細節

    免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。

    AI

    都江堰市| 海阳市| 北海市| 邢台县| 铁岭市| 柘城县| 正安县| 江北区| 马山县| 赣榆县| 秦皇岛市| 德令哈市| 开江县| 澳门| 台州市| 塔城市| 祥云县| 喀喇沁旗| 天峻县| 铅山县| 岑巩县| 鄂托克旗| 广昌县| 花垣县| 德阳市| 汽车| 朔州市| 鄢陵县| 郴州市| 滨州市| 怀集县| 荣成市| 江西省| 法库县| 五常市| 吴堡县| 登封市| 无锡市| 前郭尔| 郧西县| 平乡县|