Kafka批量消費&逐條消費詳解
更新時間:2026年02月06日 15:21:51 作者:C_Knight
文章介紹了Kafka消費者的配置參數(shù),包括批量消費和逐條消費的設(shè)置,在逐條消費模式下,消息會被分割成字符串數(shù)組,總結(jié)并提供了個人經(jīng)驗供參考
Kafka批量消費&逐條消費
消費者配置參數(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獲取實例:"+ JSON.toJSONString(factory));
return factory;
}
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory();
factory.setConsumerFactory(new DefaultKafkaConsumerFactory(defaultGoodsConsumerConfig()));
factory.setConcurrency(4);
//批量消費,如果不設(shè)置默認是單條消費
factory.setBatchListener(true);
factory.getContainerProperties().setPollTimeout(3000);
消費者監(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
如果是逐條消費,這里配置list,kafka會根據(jù)字符串中的逗號進行分割,所以碰見該現(xiàn)象不要慌,看一下批量消費的配置。
總結(jié)
以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。
相關(guān)文章
一文詳解Elasticsearch和MySQL之間的數(shù)據(jù)同步問題
Elasticsearch中的數(shù)據(jù)是來自于Mysql數(shù)據(jù)庫的,因此當數(shù)據(jù)庫中的數(shù)據(jù)進行增刪改后,Elasticsearch中的數(shù)據(jù),索引也必須跟著做出改變。本文主要來和大家探討一下Elasticsearch和MySQL之間的數(shù)據(jù)同步問題,感興趣的可以了解一下2023-04-04
Java編程實現(xiàn)時間和時間戳相互轉(zhuǎn)換實例
這篇文章主要介紹了什么是時間戳,以及Java編程實現(xiàn)時間和時間戳相互轉(zhuǎn)換實例,具有一定的參考價值,需要的朋友可以了解下。2017-09-09
springboot 如何解決static調(diào)用service為null
這篇文章主要介紹了springboot 如何解決static調(diào)用service為null的問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-06-06
基于Java編寫一個PDF與Word文件轉(zhuǎn)換工具
前段時間一直使用到word文檔轉(zhuǎn)pdf或者pdf轉(zhuǎn)word,尋思著用Java應(yīng)該是可以實現(xiàn)的,于是花了點時間寫了個文件轉(zhuǎn)換工具,感興趣的可以了解一下2023-01-01

