You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在StreamListener消费者代码中获取ConsumerRecord对象

实现方案

不需要额外配置或依赖,直接调整@StreamListener注解修饰的消费方法入参,即可拿到消息对应的分区、偏移量等元数据,根据是否开启批量消费选择对应写法即可。


单条消费场景

如果未开启批量消费,有两种实现方式:

  1. 直接接收原生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());
    }
}
  1. 轻量写法:如果不需要完整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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 12:09:18