springboot配置kafka批量消費(fèi),并發(fā)消費(fèi)方式
springboot配置kafka批量消費(fèi),并發(fā)消費(fèi)
@KafkaListener(id = "id0",groupId = "forest_fire_ql_firecard_test_info3",
topicPartitions = {@TopicPartition(topic = CommonConstants.KafkaTop.personnelCardRealTimeRecord,partitions = {"0"})},
containerFactory = "batchFactory")
public void listener0(List<String> record, Consumer<String,String> consumer){
consumer.commitSync();
try {
//業(yè)務(wù)處理
} catch (Exception e) {
log.error(e.toString());
}
}
@KafkaListener(id = "id1",groupId = "forest_fire_ql_firecard_test_info3",
topicPartitions = {@TopicPartition(topic = CommonConstants.KafkaTop.personnelCardRealTimeRecord,partitions = {"1"})},
containerFactory = "batchFactory")
public void listener1(List<String> record, Consumer<String,String> consumer){
consumer.commitSync();
try {
//業(yè)務(wù)處理
} catch (Exception e) {
log.error(e.toString());
}
}
@KafkaListener(id = "id2",groupId = "forest_fire_ql_firecard_test_info3",
topicPartitions = {@TopicPartition(topic = CommonConstants.KafkaTop.personnelCardRealTimeRecord,partitions = {"2"})},
containerFactory = "batchFactory")
public void listener2(List<String> record, Consumer<String,String> consumer){
consumer.commitSync();
try {
//業(yè)務(wù)處理
} catch (Exception e) {
log.error(e.toString());
}
}
@KafkaListener(id = "id3",groupId = "forest_fire_ql_firecard_test_info3",
topicPartitions = {@TopicPartition(topic = CommonConstants.KafkaTop.personnelCardRealTimeRecord,partitions = {"3"})},
containerFactory = "batchFactory")
public void listener3(List<String> record, Consumer<String,String> consumer){
consumer.commitSync();
try {
//業(yè)務(wù)處理
} catch (Exception e) {
log.error(e.toString());
}
}
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerContainerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;
import java.util.HashMap;
import java.util.Map;
@Slf4j
@Configuration
public class KafkaConsumerConfig {
@Value("${spring.kafka.bootstrap-servers}")
private String bootstrapServersConfig;
public Map<String,Object> consumerConfigs(){
Map<String,Object> props = new HashMap<>();
props.put(ConsumerConfig.GROUP_ID_CONFIG, "forest_fire_ql_firecard_test_info3");
log.info("bootstrapServersConfig:自定義配置="+ bootstrapServersConfig);
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServersConfig);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG,3);
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,false);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"earliest");
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG,"100");
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG,"20000");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
return props;
}
@Bean
public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, Object>> batchFactory(KafkaProperties properties) {
//Map<String, Object> consumerProperties = properties.buildConsumerProperties();
ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigs()));
//并發(fā)數(shù)量
factory.setConcurrency(3);
//開啟批量監(jiān)聽,消費(fèi)
factory.setBatchListener(true);
//factory.set
return factory;
}
}按照以上配置內(nèi)容即可,可以達(dá)到kafka批量消費(fèi)的能力。
但是,要特別需要注意的一個(gè)點(diǎn)是:
- 并發(fā)量根據(jù)實(shí)際的分區(qū)數(shù)量決定
- 必須小于等于分區(qū)數(shù)
- 否則會(huì)有線程一直處于空閑狀態(tài)
下面是創(chuàng)建4個(gè)分區(qū)的命令寫法
bin/kafka-topics.sh --create --bootstrap-server localhost:9092 --topic personnel_card_real_time_recordinfo --partitions 4 --replication-factor 1
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
- Springboot項(xiàng)目消費(fèi)Kafka數(shù)據(jù)的方法
- SpringBoot整合Kafka完成生產(chǎn)消費(fèi)的方案
- springboot如何開啟和關(guān)閉kafka消費(fèi)
- spring-kafka使消費(fèi)者動(dòng)態(tài)訂閱新增的topic問題
- springboot集成kafka消費(fèi)手動(dòng)啟動(dòng)停止操作
- Springboot集成Kafka進(jìn)行批量消費(fèi)及踩坑點(diǎn)
- spring 整合kafka監(jiān)聽消費(fèi)的配置過程
- 淺談spring-kafka消費(fèi)異常處理
相關(guān)文章
jdk配置完之后java?-version還是默認(rèn)的jdk版本問題解決過程
在CentOS7中配置JDK后,java?-version仍顯示默認(rèn)JDK,因?yàn)橄到y(tǒng)默認(rèn)的/usr/bin/java軟鏈接優(yōu)先級(jí)高于PATH環(huán)境變量,這篇文章主要介紹了jdk配置完之后java?-version還是默認(rèn)的jdk版本問題的解決過程,需要的朋友可以參考下2026-01-01
JAXB解析xml轉(zhuǎn)換成類的實(shí)現(xiàn)方式
本文主要介紹了如何使用JAXB將XML配置項(xiàng)轉(zhuǎn)換為Java類,JAXB提供了多種注解,如@XmlRootElement、@XmlElement、@XmlElementWrapper、@XmlAttribute等,可以方便地將XML元素映射為Java對(duì)象,并且可以控制生成的XML結(jié)構(gòu),同時(shí),文章也提到了一些需要注意的問題2025-11-11
詳解java創(chuàng)建一個(gè)女朋友類(對(duì)象啥的new一個(gè)就是)==建造者模式,一鍵重寫
這篇文章主要介紹了java建造者模式一鍵重寫,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2019-04-04
java使用POI批量導(dǎo)入excel數(shù)據(jù)的方法
這篇文章主要為大家詳細(xì)介紹了java使用POI批量導(dǎo)入excel數(shù)據(jù)的方法,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2017-07-07
Java實(shí)戰(zhàn)之酒店人事管理系統(tǒng)的實(shí)現(xiàn)
這篇文章主要介紹了如何用Java實(shí)現(xiàn)酒店人事管理系統(tǒng),文中采用的技術(shù)有:JSP、Spring、SpringMVC、MyBatis等,感興趣的小伙伴可以學(xué)習(xí)一下2022-03-03
SpringBoot集成Session的實(shí)現(xiàn)示例
Session是一個(gè)在Web開發(fā)中常用的概念,它表示服務(wù)器和客戶端之間的一種狀態(tài)管理機(jī)制,用于跟蹤用戶在網(wǎng)站或應(yīng)用程序中的狀態(tài)和數(shù)據(jù),本文主要介紹了SpringBoot集成Session的實(shí)現(xiàn)示例,感興趣的可以了解一下2023-09-09

