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

Spring?Boot?中使用@KafkaListener并發(fā)批量接收消息的完整代碼

 更新時(shí)間:2023年02月20日 09:32:59   作者:russle  
kakfa是我們?cè)陧?xiàng)目開(kāi)發(fā)中經(jīng)常使用的消息中間件。由于它的寫(xiě)性能非常高,因此,經(jīng)常會(huì)碰到讀取Kafka消息隊(duì)列時(shí)擁堵的情況,這篇文章主要介紹了Spring?Boot?中使用@KafkaListener并發(fā)批量接收消息,需要的朋友可以參考下

kakfa是我們?cè)陧?xiàng)目開(kāi)發(fā)中經(jīng)常使用的消息中間件。由于它的寫(xiě)性能非常高,因此,經(jīng)常會(huì)碰到讀取Kafka消息隊(duì)列時(shí)擁堵的情況。遇到這種情況時(shí),有時(shí)我們不能直接清理整個(gè)topic,因?yàn)檫€有別的服務(wù)正在使用該topic。因此只能額外啟動(dòng)一個(gè)相同名稱(chēng)的consumer-group來(lái)加快消息消費(fèi)(如果該topic只有一個(gè)分區(qū),再啟動(dòng)一個(gè)新的消費(fèi)者,沒(méi)有作用)。

完整的代碼在這里,歡迎加星號(hào)、fork。

官方文檔在https://docs.spring.io/spring-kafka/reference/html/_reference.html

###第一步,并發(fā)消費(fèi)###
先看代碼,重點(diǎn)是這我們使用的是ConcurrentKafkaListenerContainerFactory并且設(shè)置了factory.setConcurrency(4); (我的topic有4個(gè)分區(qū),為了加快消費(fèi)將并發(fā)設(shè)置為4,也就是有4個(gè)KafkaMessageListenerContainer)

    @Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setConcurrency(4);
        factory.setBatchListener(true);
        factory.getContainerProperties().setPollTimeout(3000);
        return factory;
    }

注意也可以直接在application.properties中添加spring.kafka.listener.concurrency=3,然后使用@KafkaListener并發(fā)消費(fèi)。

###第二步,批量消費(fèi)###
然后是批量消費(fèi)。重點(diǎn)是factory.setBatchListener(true);
以及 propsMap.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 50);
一個(gè)設(shè)啟用批量消費(fèi),一個(gè)設(shè)置批量消費(fèi)每次最多消費(fèi)多少條消息記錄。

重點(diǎn)說(shuō)明一下,我們?cè)O(shè)置的ConsumerConfig.MAX_POLL_RECORDS_CONFIG是50,并不是說(shuō)如果沒(méi)有達(dá)到50條消息,我們就一直等待。官方的解釋是"The maximum number of records returned in a single call to poll().", 也就是50表示的是一次poll最多返回的記錄數(shù)。

從啟動(dòng)日志中可以看到還有個(gè) max.poll.interval.ms = 300000, 也就說(shuō)每間隔max.poll.interval.ms我們就調(diào)用一次poll。每次poll最多返回50條記錄。

max.poll.interval.ms官方解釋是"The maximum delay between invocations of poll() when using consumer group management. This places an upper bound on the amount of time that the consumer can be idle before fetching more records. If poll() is not called before expiration of this timeout, then the consumer is considered failed and the group will rebalance in order to reassign the partitions to another member. ";

    @Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setConcurrency(4);
        factory.setBatchListener(true);
        factory.getContainerProperties().setPollTimeout(3000);
        return factory;
    }

   @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> propsMap = new HashMap<>();
        propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, propsConfig.getBroker());
        propsMap.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, propsConfig.getEnableAutoCommit());
        propsMap.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "100");
        propsMap.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
        propsMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(ConsumerConfig.GROUP_ID_CONFIG, propsConfig.getGroupId());
        propsMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, propsConfig.getAutoOffsetReset());
        propsMap.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 50);
        return propsMap;
    }

啟動(dòng)日志截圖

這里寫(xiě)圖片描述

關(guān)于max.poll.records和max.poll.interval.ms官方解釋截圖:

這里寫(xiě)圖片描述

###第三步,分區(qū)消費(fèi)###
對(duì)于只有一個(gè)分區(qū)的topic,不需要分區(qū)消費(fèi),因?yàn)闆](méi)有意義。下面的例子是針對(duì)有2個(gè)分區(qū)的情況(我的完整代碼中有4個(gè)listenPartitionX方法,我的topic設(shè)置了4個(gè)分區(qū)),讀者可以根據(jù)自己的情況進(jìn)行調(diào)整。

public class MyListener {
    private static final String TPOIC = "topic02";

    @KafkaListener(id = "id0", topicPartitions = { @TopicPartition(topic = TPOIC, partitions = { "0" }) })
    public void listenPartition0(List<ConsumerRecord<?, ?>> records) {
        log.info("Id0 Listener, Thread ID: " + Thread.currentThread().getId());
        log.info("Id0 records size " +  records.size());

        for (ConsumerRecord<?, ?> record : records) {
            Optional<?> kafkaMessage = Optional.ofNullable(record.value());
            log.info("Received: " + record);
            if (kafkaMessage.isPresent()) {
                Object message = record.value();
                String topic = record.topic();
                log.info("p0 Received message={}",  message);
            }
        }
    }

    @KafkaListener(id = "id1", topicPartitions = { @TopicPartition(topic = TPOIC, partitions = { "1" }) })
    public void listenPartition1(List<ConsumerRecord<?, ?>> records) {
        log.info("Id1 Listener, Thread ID: " + Thread.currentThread().getId());
        log.info("Id1 records size " +  records.size());

        for (ConsumerRecord<?, ?> record : records) {
            Optional<?> kafkaMessage = Optional.ofNullable(record.value());
            log.info("Received: " + record);
            if (kafkaMessage.isPresent()) {
                Object message = record.value();
                String topic = record.topic();
                log.info("p1 Received message={}",  message);
            }
        }
}

關(guān)于分區(qū)和消費(fèi)者關(guān)系,后面會(huì)補(bǔ)充,先摘錄如下:
If, say, 6 TopicPartition s are provided and the concurrency is 3; each container will get 2 partitions. For 5 TopicPartition s, 2 containers will get 2 partitions and the third will get 1. If the concurrency is greater than the number of TopicPartitions, the concurrency will be adjusted down such that each container will get one partition.

最后,總結(jié),如果我們的topic有多個(gè)分區(qū),經(jīng)過(guò)以上步驟可以很好的加快消息消費(fèi)。如果只有一個(gè)分區(qū),因?yàn)橐呀?jīng)有一個(gè)同名group id在消費(fèi)了,新啟動(dòng)的一個(gè)基本上沒(méi)有作用(本人測(cè)試結(jié)果)。

具體代碼在這里,歡迎加星號(hào),fork。

到此這篇關(guān)于Spring Boot 中使用@KafkaListener并發(fā)批量接收消息的文章就介紹到這了,更多相關(guān)Spring Boot 使用@KafkaListener并發(fā)批量接收消息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java中常用的四種引用類(lèi)型詳解

    Java中常用的四種引用類(lèi)型詳解

    Java中常用的四種引用類(lèi)型,分別為,強(qiáng)引用、軟引用、弱引用以及虛引用,這篇文章主要為大家介紹了這四種引用的用法,需要的可以參考一下
    2023-06-06
  • Spring常用數(shù)據(jù)源的xml配置詳解

    Spring常用數(shù)據(jù)源的xml配置詳解

    這篇文章主要介紹了Spring常用數(shù)據(jù)源的xml配置詳解,數(shù)據(jù)源是連接到數(shù)據(jù)庫(kù)的一類(lèi)路徑,它包含了訪問(wèn)數(shù)據(jù)庫(kù)的信息(地址、用戶名、密碼),數(shù)據(jù)源就像是排水管道,需要的朋友可以參考下
    2023-07-07
  • spring?Bean創(chuàng)建的完整過(guò)程記錄

    spring?Bean創(chuàng)建的完整過(guò)程記錄

    這篇文章主要給大家介紹了關(guān)于Spring中Bean實(shí)例創(chuàng)建的相關(guān)資料,文中通過(guò)實(shí)例代碼和圖文介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2022-01-01
  • java環(huán)境變量path和classpath的配置

    java環(huán)境變量path和classpath的配置

    這篇文章主要為大家詳細(xì)介紹了java系統(tǒng)環(huán)境變量path和classpath的配置過(guò)程,感興趣的小伙伴們可以參考一下
    2016-07-07
  • Java pdf和jpg互轉(zhuǎn)案例

    Java pdf和jpg互轉(zhuǎn)案例

    這篇文章主要介紹了Java pdf和jpg互轉(zhuǎn)案例,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-09-09
  • Java多線程的臨界資源問(wèn)題解決方案

    Java多線程的臨界資源問(wèn)題解決方案

    這篇文章主要介紹了Java多線程的臨界資源問(wèn)題解決方案,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-02-02
  • 基于mybatis高級(jí)映射多對(duì)多查詢的實(shí)現(xiàn)

    基于mybatis高級(jí)映射多對(duì)多查詢的實(shí)現(xiàn)

    下面小編就為大家?guī)?lái)一篇基于mybatis高級(jí)映射多對(duì)多查詢的實(shí)現(xiàn)。小編覺(jué)得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2017-10-10
  • java基于socket傳輸zip文件功能示例

    java基于socket傳輸zip文件功能示例

    這篇文章主要介紹了java基于socket傳輸zip文件功能,結(jié)合實(shí)例形式分析了java使用socket進(jìn)行文件傳輸?shù)木唧w操作步驟與服務(wù)器端、客戶端相關(guān)實(shí)現(xiàn)技巧,需要的朋友可以參考下
    2017-07-07
  • SpringCloud配置客戶端ConfigClient接入服務(wù)端

    SpringCloud配置客戶端ConfigClient接入服務(wù)端

    這篇文章主要為大家介紹了SpringCloud配置客戶端ConfigClient接入服務(wù)端,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2023-08-08
  • 淺談異常結(jié)構(gòu)圖、編譯期異常和運(yùn)行期異常的區(qū)別

    淺談異常結(jié)構(gòu)圖、編譯期異常和運(yùn)行期異常的區(qū)別

    下面小編就為大家?guī)?lái)一篇淺談異常結(jié)構(gòu)圖、編譯期異常和運(yùn)行期異常的區(qū)別。小編覺(jué)得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2016-09-09

最新評(píng)論

亚东县| 建阳市| 政和县| 肇庆市| 仁寿县| 苏尼特左旗| 四会市| 通化县| 西青区| 丽水市| 钦州市| 那坡县| 宿松县| 巴林左旗| 丰宁| 盐边县| 高雄县| 衢州市| 永定县| 富阳市| 弋阳县| 岳阳县| 萨迦县| 西藏| 江源县| 宁武县| 新闻| 中超| 中方县| 金湖县| 万荣县| 赫章县| 通山县| 明溪县| 永仁县| 浮山县| 阳春市| 应城市| 资阳市| 承德市| 尚志市|