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

Kafka批量消費(fèi)&逐條消費(fèi)詳解

 更新時(shí)間:2026年02月06日 15:21:51   作者:C_Knight  
文章介紹了Kafka消費(fèi)者的配置參數(shù),包括批量消費(fèi)和逐條消費(fèi)的設(shè)置,在逐條消費(fèi)模式下,消息會(huì)被分割成字符串?dāng)?shù)組,總結(jié)并提供了個(gè)人經(jīng)驗(yàn)供參考

Kafka批量消費(fèi)&逐條消費(fèi)

消費(fèi)者配置參數(shù)

    private Map<String, Object> defaultGoodsConsumerConfig() {
        Map<String, Object> props = Maps.newHashMap();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "ip:port");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
        props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "100");
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "50");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer
");
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer
");
        props.put("listener.type", "batch");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "modify-group");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, “SASL_PLAINTEXT”);
        props.put(SaslConfigs.SASL_MECHANISM, defaultKafkaProperties.getSaslMechanism());
        props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username="###" password="###";
");

        return props;
    }

    @Bean(name = "defaultListenerContainerFactory")
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> defaultListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory();
        factory.setConsumerFactory(new DefaultKafkaConsumerFactory(defaultGoodsConsumerConfig()));
        factory.setConcurrency(4);
        factory.setBatchListener(true);
        factory.getContainerProperties().setPollTimeout(3000);
        log.info("KafkaDefaultConsumer factory獲取實(shí)例:"+ JSON.toJSONString(factory));
        return factory;
    }
       ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory();
       factory.setConsumerFactory(new DefaultKafkaConsumerFactory(defaultGoodsConsumerConfig()));
       factory.setConcurrency(4);
       //批量消費(fèi),如果不設(shè)置默認(rèn)是單條消費(fèi)
       factory.setBatchListener(true);
       factory.getContainerProperties().setPollTimeout(3000);

消費(fèi)者監(jiān)聽消息

    /**
     * 監(jiān)聽goods變更消息
     */
    @KafkaListener(id="sync-modify-goods", topics = "${kafka.sync.goods.topic}", concurrency = "4", containerFactory = "defaultListenerContainerFactory")
    public void updateListener(List<ConsumerRecord<String, String>> records){
        for (ConsumerRecord<String, String> msg:records) {
            GoodsChangeMsg changeMsg = null;
            try {
                changeMsg = JSONObject.parseObject(msg.value(), GoodsChangeMsg.class);
                syncGoodsProcessor.handle(changeMsg);
            }catch (Exception exception) {
                log.error("解析失敗{}", msg, exception);
            }
        }
    }

List<ConsumerRecord<String, String>> records可以是String[] message

如果是逐條消費(fèi),這里配置list,kafka會(huì)根據(jù)字符串中的逗號(hào)進(jìn)行分割,所以碰見該現(xiàn)象不要慌,看一下批量消費(fèi)的配置。

總結(jié)

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

相關(guān)文章

  • Java實(shí)現(xiàn)WebSocket四個(gè)步驟

    Java實(shí)現(xiàn)WebSocket四個(gè)步驟

    這篇文章主要為大家介紹了Java實(shí)現(xiàn)WebSocket的方法實(shí)例,只需要簡(jiǎn)單四個(gè)步驟,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2024-01-01
  • 一文詳解Elasticsearch和MySQL之間的數(shù)據(jù)同步問題

    一文詳解Elasticsearch和MySQL之間的數(shù)據(jù)同步問題

    Elasticsearch中的數(shù)據(jù)是來(lái)自于Mysql數(shù)據(jù)庫(kù)的,因此當(dāng)數(shù)據(jù)庫(kù)中的數(shù)據(jù)進(jìn)行增刪改后,Elasticsearch中的數(shù)據(jù),索引也必須跟著做出改變。本文主要來(lái)和大家探討一下Elasticsearch和MySQL之間的數(shù)據(jù)同步問題,感興趣的可以了解一下
    2023-04-04
  • 詳細(xì)了解JAVA NIO之Buffer(緩沖區(qū))

    詳細(xì)了解JAVA NIO之Buffer(緩沖區(qū))

    這篇文章主要介紹了JAVA NIO之Buffer(緩沖區(qū))的相關(guān)資料,文中講解非常細(xì)致,幫助大家更好的學(xué)習(xí)JAVA NIO,感興趣的朋友可以了解下
    2020-07-07
  • SpringBoot整合EasyExcel的完整過(guò)程記錄

    SpringBoot整合EasyExcel的完整過(guò)程記錄

    easyexcel是阿里巴巴旗下開源項(xiàng)目,主要用于Excel文件的導(dǎo)入和導(dǎo)出處理,下面這篇文章主要給大家介紹了關(guān)于SpringBoot整合EasyExcel的完整過(guò)程,需要的朋友可以參考下
    2021-12-12
  • Java8?Stream流的常用方法匯總

    Java8?Stream流的常用方法匯總

    Java8?API添加了一個(gè)新的抽象稱為流Stream,可以讓你以一種聲明的方式處理數(shù)據(jù),下面這篇文章主要給大家介紹了關(guān)于Java8?Stream流的常用方法,文中通過(guò)示例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-07-07
  • Java編程實(shí)現(xiàn)時(shí)間和時(shí)間戳相互轉(zhuǎn)換實(shí)例

    Java編程實(shí)現(xiàn)時(shí)間和時(shí)間戳相互轉(zhuǎn)換實(shí)例

    這篇文章主要介紹了什么是時(shí)間戳,以及Java編程實(shí)現(xiàn)時(shí)間和時(shí)間戳相互轉(zhuǎn)換實(shí)例,具有一定的參考價(jià)值,需要的朋友可以了解下。
    2017-09-09
  • 聊聊java中引用數(shù)據(jù)類型有哪些

    聊聊java中引用數(shù)據(jù)類型有哪些

    這篇文章主要介紹了java中引用數(shù)據(jù)類型有哪些,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-10-10
  • SpringBoot實(shí)現(xiàn)線程池

    SpringBoot實(shí)現(xiàn)線程池

    現(xiàn)在由于系統(tǒng)越來(lái)越復(fù)雜,導(dǎo)致很多接口速度變慢,這時(shí)候就會(huì)想到可以利用線程池來(lái)處理一些耗時(shí)并不影響系統(tǒng)的操作。本文就介紹了SpringBoot線程池的使用,感興趣的可以了解一下
    2021-06-06
  • springboot 如何解決static調(diào)用service為null

    springboot 如何解決static調(diào)用service為null

    這篇文章主要介紹了springboot 如何解決static調(diào)用service為null的問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • 基于Java編寫一個(gè)PDF與Word文件轉(zhuǎn)換工具

    基于Java編寫一個(gè)PDF與Word文件轉(zhuǎn)換工具

    前段時(shí)間一直使用到word文檔轉(zhuǎn)pdf或者pdf轉(zhuǎn)word,尋思著用Java應(yīng)該是可以實(shí)現(xiàn)的,于是花了點(diǎn)時(shí)間寫了個(gè)文件轉(zhuǎn)換工具,感興趣的可以了解一下
    2023-01-01

最新評(píng)論

抚松县| 安阳市| 大埔县| 日土县| 垫江县| 茂名市| 博乐市| 祥云县| 百色市| 香格里拉县| 双鸭山市| 望江县| 五华县| 龙井市| 赤峰市| 玛曲县| 宁陕县| 永平县| 喀喇| 略阳县| 阿合奇县| 河北区| 封丘县| 新绛县| 昌乐县| 鹤庆县| 合作市| 鹤峰县| 大兴区| 囊谦县| 莱西市| 耒阳市| 平远县| 错那县| 平果县| 湘乡市| 奉贤区| 庄河市| 邻水| 政和县| 长汀县|