如何在StreamListener消费者代码中获取ConsumerRecord对象
实现方案
不需要额外配置或依赖,直接调整@StreamListener注解修饰的消费方法入参,即可拿到消息对应的分区、偏移量等元数据,根据是否开启批量消费选择对应写法即可。
单条消费场景
如果未开启批量消费,有两种实现方式:
- 直接接收原生
ConsumerRecord对象,框架会自动注入完整的消息元数据:
@StreamListener(ConsumerConstants.COMMUNITY_IN) public void handleCommFeedConsumer( ConsumerRecord<String, Account> record, @Header(KafkaHeaders.CONSUMER) Consumer<?, ?> consumer, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) { Account communityFeed = record.value(); try { AccountClient.signIn(AccountIn.builder() .Id(communityFeed.getId()) .build()); log.debug("Calling Client for Id : {}", communityFeed.getId()); acknowledgment.acknowledge(); } catch (RuntimeException ex) { log.error("消息处理失败,即将重置偏移量,topic:{}, partition:{}, offset:{}", record.topic(), record.partition(), record.offset(), ex); TopicPartition tp = new TopicPartition(record.topic(), record.partition()); consumer.seek(tp, record.offset()); } }
- 轻量写法:如果不需要完整ConsumerRecord对象,直接通过消息头注入所需元数据即可:
@StreamListener(ConsumerConstants.COMMUNITY_IN) public void handleCommFeedConsumer( @Payload Account communityFeed, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition, @Header(KafkaHeaders.OFFSET) long offset, @Header(KafkaHeaders.CONSUMER) Consumer<?, ?> consumer, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) { try { AccountClient.signIn(AccountIn.builder() .Id(communityFeed.getId()) .build()); acknowledgment.acknowledge(); } catch (RuntimeException ex) { consumer.seek(new TopicPartition(topic, partition), offset); } }
批量消费场景
从现有代码的forEach逻辑判断,大概率开启了批量消费,这种场景下消息头只会存储整批最后一条消息的元数据,必须直接接收List<ConsumerRecord>类型的入参,才能拿到每条消息的独立元数据,参考实现如下:
@StreamListener(ConsumerConstants.COMMUNITY_IN) public void handleCommFeedConsumer( List<ConsumerRecord<String, Account>> consumerRecords, @Header(KafkaHeaders.CONSUMER) Consumer<?, ?> consumer, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) { Map<TopicPartition, Long> seekTargetMap = new HashMap<>(); boolean processSuccess = true; for (ConsumerRecord<String, Account> record : consumerRecords) { Account communityFeed = record.value(); try { AccountClient.signIn(AccountIn.builder() .Id(communityFeed.getId()) .build()); log.debug("Calling Client for Id : {}", communityFeed.getId()); } catch (RuntimeException ex) { log.error("消息处理失败,topic:{}, partition:{}, offset:{}", record.topic(), record.partition(), record.offset(), ex); TopicPartition tp = new TopicPartition(record.topic(), record.partition()); // 同一分区只保留最小的失败偏移量,避免漏消费 seekTargetMap.merge(tp, record.offset(), Math::min); processSuccess = false; // 遇到失败消息直接终止后续处理,等待重消费 break; } } if (processSuccess) { acknowledgment.acknowledge(); } else { seekTargetMap.forEach(consumer::seek); } }
注意事项:
- 批量消费必须保证配置项
spring.cloud.stream.kafka.bindings.atcommnity.consumer.batch-mode=true,否则框架无法注入List<ConsumerRecord>类型参数- 批量场景禁止在循环内部调用
acknowledge(),该方法会提交当前批次最大的消息偏移量,会导致未处理成功的消息被误提交- 执行seek操作时直接传入失败消息的原始offset即可,不需要做+1处理,否则会跳过当前失败消息
- 开启手动提交偏移量时,只要不调用
acknowledge(),本次拉取的批次偏移量就不会被提交,配合seek可以精准回到失败消息的位置实现重消费
内容的提问来源于stack exchange,提问作者nitin737
相关产品推荐
相关产品推荐

