news 2026/9/4 20:58:54

Kafka批量消费实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Kafka批量消费实现

批量消费指的是一次性拉取一批消息,然后批量处理
依赖spring-kafka

<dependency><groupId>org.springframework.kafka</groupId><artifactId>spring-kafka</artifactId><version>2.2.4.RELEASE</version></dependency>

配置消费者工厂

@ConfigurationpublicclassKafkaConsumerConfig{@BeanpublicConsumerFactory<String,String>consumerFactory(){Map<String,Object>props=newHashMap<>();props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092");props.put(ConsumerConfig.GROUP_ID_CONFIG,"batch-consumer-group");props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class);props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class);// 批量消费配置props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG,100);// 每次最多拉取100条props.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG,10240);// 至少10KB才返回returnnewDefaultKafkaConsumerFactory<>(props);}@BeanpublicConcurrentKafkaListenerContainerFactory<String,String>batchFactory(ConsumerFactory<String,String>consumerFactory){ConcurrentKafkaListenerContainerFactory<String,String>factory=newConcurrentKafkaListenerContainerFactory<>();factory.setConsumerFactory(consumerFactory);factory.setBatchListener(true);// 启用批量监听// 手动提交模式factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);returnfactory;}}

实现批量消费监听器

@ServicepublicclassOrderBatchConsumer{@KafkaListener(topics="order-topic",containerFactory="batchFactory")publicvoidconsumeBatch(List<ConsumerRecord<String,String>>records,Acknowledgmentack){log.info("收到批量消息: {} 条",records.size());try{// 业务处理for(ConsumerRecord<String,String>record:records){processOrder(record.value());}// 全部成功才提交offsetack.acknowledge();log.info("批次处理成功,已提交offset");}catch(Exceptione){log.error("批量处理失败",e);// 不提交offset,等待重新投递}}}

批量消费可靠性保障
不能在finally中提交offset,因为不管消费是否成功都会提交offset

@KafkaListener(topics="topic",containerFactory="batchFactory")publicvoidconsume(List<ConsumerRecord<String,String>>records,Acknowledgmentack){try{processBatch(records);}finally{ack.acknowledge();// 无论成败都提交,会丢消息!}}

不能设置自动提交,自动提交模式下,批量消费有可能会丢消息。要改为手动提交

props.put("enable.auto.commit","true");// 自动提交// 处理过程中失败,但offset已自动提交,消息丢失

需要用到线程池的时候,需要等到全部消息处理成功才能提交。使用CompletionService实现

@ServicepublicclassReliableBatchConsumer{privatefinalExecutorServiceexecutor=Executors.newFixedThreadPool(10);@KafkaListener(topics="payment-topic",containerFactory="batchFactory")publicvoidconsumeBatch(List<ConsumerRecord<String,String>>records,Acknowledgmentack){CompletionService<Boolean>completionService=newExecutorCompletionService<>(executor);List<Future<Boolean>>futures=newArrayList<>();// 1. 提交所有任务到线程池并发处理for(ConsumerRecord<String,String>record:records){Callable<Boolean>task=()->{try{processPayment(record.value());returntrue;}catch(Exceptione){log.error("支付处理失败: {}",record.value(),e);returnfalse;}};futures.add(completionService.submit(task));}// 2. 等待所有任务完成并检查结果booleanallSuccess=true;try{for(inti=0;i<records.size();i++){Future<Boolean>future=completionService.take();if(!future.get()){allSuccess=false;break;// 发现失败立即终止}}}catch(Exceptione){allSuccess=false;log.error("任务执行异常",e);}// 3. 全部成功才提交offsetif(allSuccess){ack.acknowledge();log.info("批次全部成功,已提交offset");}else{log.warn("批次中有失败消息,不提交offset,等待重新投递");// 重新投递会导致重复消费,需要业务保证幂等性}}}

当有失败消息的时候,不提交offset,等待消息重投或者重新拉取,但是会有消息重复的情况,业务上要做好幂等。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/3 5:50:07

PHP日志格式设计陷阱:80%开发者忽略的3个致命问题

第一章&#xff1a;PHP日志格式设计陷阱&#xff1a;80%开发者忽略的3个致命问题非结构化日志导致排查困难 许多PHP项目仍采用简单的 error_log() 输出文本日志&#xff0c;缺乏统一结构。这使得在系统出错时难以快速定位关键信息。// 错误示例&#xff1a;非结构化输出 error_…

作者头像 李华
网站建设 2026/9/3 2:01:37

【PHP 8.7扩展开发终极指南】:手把手教你从零编写高性能C扩展

第一章&#xff1a;PHP 8.7扩展开发概述PHP 扩展开发是深入理解 PHP 内核机制的重要途径&#xff0c;尤其在 PHP 8.7 即将发布的背景下&#xff0c;扩展开发能力对于性能优化、功能定制和底层集成具有重要意义。通过编写 C 语言实现的扩展&#xff0c;开发者可以直接与 Zend 引…

作者头像 李华
网站建设 2026/9/3 2:55:27

语音合成中的性别转换技术:男声转女声自然度实测

语音合成中的性别转换技术&#xff1a;男声转女声自然度实测 在虚拟主播越来越像真人、AI客服开始带情绪说话的今天&#xff0c;我们早已不再满足于“能出声”的TTS系统。越来越多的应用场景提出了更细腻的需求——比如让一个低沉的男声&#xff0c;自然地变成温柔或干练的女声…

作者头像 李华
网站建设 2026/9/3 1:23:34

javascript postMessage跨域通信调用外部GLM-TTS

JavaScript postMessage 实现跨域调用本地 GLM-TTS 语音合成 在如今的前端架构中&#xff0c;AI能力正越来越多地被集成到Web应用中——从图像生成、语音识别&#xff0c;到更复杂的文本到语音合成&#xff08;TTS&#xff09;。然而&#xff0c;这些模型往往体积庞大、依赖复杂…

作者头像 李华
网站建设 2026/9/2 22:36:49

仅限高级开发者掌握的PHP跨域防护技巧(限时公开)

第一章&#xff1a;PHP跨域安全策略的核心认知在现代Web开发中&#xff0c;前后端分离架构已成为主流&#xff0c;PHP作为后端服务常需处理来自不同源的前端请求。浏览器基于同源策略&#xff08;Same-Origin Policy&#xff09;限制跨域资源访问&#xff0c;以防止恶意攻击。因…

作者头像 李华
网站建设 2026/9/3 0:23:53

打造高安全性视频平台:PHP+FFmpeg+AES加密全流程详解(独家方案)

第一章&#xff1a;高安全性视频平台架构概述在构建现代视频服务平台时&#xff0c;安全性已成为核心设计原则之一。高安全性视频平台不仅需保障用户数据的机密性与完整性&#xff0c;还需防范非法访问、内容盗取及服务中断等风险。为此&#xff0c;系统架构从传输层到应用层均…

作者头像 李华