SpringBoot中KafkaConsumer未处理全部请求 部分数据未存入MongoDB
问题根因定位
- 批量监听配置与消费方法参数不匹配:你在
userKafkaListenerFactory中设置了factory.setBatchListener(true)开启批量监听模式,该模式下Kafka会一次性投递一批消息到消费方法,要求方法入参必须是List<UserKafkaDTO>类型。你当前方法入参是单个UserKafkaDTO,只会处理批量中的第一条消息,剩余消息直接被丢弃,这是投递5条仅存3条的核心原因。 - 手动偏移量提交未实现:你配置了
ENABLE_AUTO_COMMIT_CONFIG = "false"关闭了自动偏移量提交,但消费逻辑中没有任何手动提交偏移量的代码,就算消息消费成功,偏移量也不会更新,若服务重启会重复消费未提交偏移量的消息;如果消费过程抛出异常未捕获,当前批次偏移量不会提交,也会导致消息看似丢失。 - 异常处理缺失:
invokeAPItoSaveRecordTOMongoDB调用REST接口存MongoDB的过程没有异常捕获逻辑,一旦接口超时、MongoDB写入失败抛出异常,当前消息不会走任何重试/降级逻辑,直接丢失。
修复方案
1. 修正批量监听的消费方法入参
把消费方法的入参改为List<UserKafkaDTO>,遍历处理每一条消息:
@KafkaListener(topics = { "topicName" }, containerFactory = "userKafkaListenerFactory",autoStartup = "${listen.auto.start}") public void consumeJson(List<UserKafkaDTO> userKafkaDTOList) { for (UserKafkaDTO dto : userKafkaDTOList) { invokeAPItoSaveRecordTOMongoDB(dto); } }
如果不需要批量监听,直接把factory.setBatchListener(true)删除或者改为false,即可使用单个参数的消费方法。
2. 补充偏移量提交逻辑
可选择两种实现方式:
- 方式一:开启自动提交,把
ENABLE_AUTO_COMMIT_CONFIG改为"true",同时配置AUTO_COMMIT_INTERVAL_MS_CONFIG设置自动提交间隔,适合对消息一致性要求不高的场景。 - 方式二:保留手动提交,在
@KafkaListener方法中加入Acknowledgment参数,处理完所有消息后手动提交:
@KafkaListener(topics = { "topicName" }, containerFactory = "userKafkaListenerFactory",autoStartup = "${listen.auto.start}") public void consumeJson(List<UserKafkaDTO> userKafkaDTOList, Acknowledgment ack) { try { for (UserKafkaDTO dto : userKafkaDTOList) { invokeAPItoSaveRecordTOMongoDB(dto); } // 全部处理成功后提交偏移量 ack.acknowledge(); } catch (Exception e) { // 自行实现失败处理逻辑:重试、写入死信队列等 e.printStackTrace(); } }
同时需要在userKafkaListenerFactory配置中添加手动ack模式:
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
3. 增加异常重试与降级机制
在invokeAPItoSaveRecordTOMongoDB外层添加try-catch捕获异常,写入失败的消息可以先存入本地死信表,后续定时补偿重试,避免消息直接丢失。
内容的提问来源于stack exchange,提问作者Amulya M
相关产品推荐
相关产品推荐

