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

Spring-Kafka异步手动提交:记录Offset提交耗时

问题

我有一个基于Spring-Kafka的批量消息监听器,处理完成后手动确认消息,代码如下:

@KafkaListener(topics = my_topic)
public void consume(@Payload List<String> messages,
                    @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) {
    for (String message: messages) {
         // 处理每条消息
    }
    // 理想情况下,计时器应在此处记录开始时间
    acknowledgment.acknowledge();
}

配置类中包含提交回调逻辑,我希望记录从调用acknowledge()方法到收到回调之间的耗时,但acknowledge()无法传递任何值到回调中:

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {

    log.info("创建Kafka监听器工厂");
    ConcurrentKafkaListenerContainerFactory<String, String> factory
            = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setBatchListener(true);
    factory.getContainerProperties().setSyncCommits(false);     
    factory.getContainerProperties()
            .setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
    // 更多配置
    factory.getContainerProperties().setCommitCallback((map, ex) -> {
        // 理想情况下,应在此处记录结束时间并发布计时器指标
        if (ex == null) {
            log.info("提交成功:{}", map);
        } else {
            log.error("提交失败:{},异常信息:{}", map, ex.getMessage());
        }
    });
    return factory;
}

如何实现记录这段耗时的需求?

解决方案

可以通过缓存提交开始时间+关联偏移量的方式实现,核心思路是在调用acknowledge()前记录时间,通过提交回调中的偏移量映射关系匹配对应的开始时间,计算耗时。

步骤1:添加全局缓存与监听器参数

在监听器所在类中定义线程安全的缓存,存储偏移量对应的提交开始时间;同时在监听器方法中获取批量消息对应的分区和偏移量:

import org.apache.kafka.common.TopicPartition;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

@Component
public class KafkaBatchListener {

    // 缓存:key为TopicPartition+偏移量的字符串,value为提交开始时间戳
    private final Map<String, Long> commitStartTimeCache = new ConcurrentHashMap<>();

    @KafkaListener(topics = "my_topic")
    public void consume(@Payload List<String> messages,
                        @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment,
                        @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions,
                        @Header(KafkaHeaders.OFFSET) List<Long> offsets) {
        for (String message : messages) {
            // 处理每条消息
        }

        // 记录每个偏移量的提交开始时间
        for (int i = 0; i < partitions.size(); i++) {
            String key = "my_topic-" + partitions.get(i) + "-" + offsets.get(i);
            commitStartTimeCache.put(key, System.currentTimeMillis());
        }

        acknowledgment.acknowledge();
    }
}

步骤2:修改提交回调计算耗时

在配置类的提交回调中,遍历提交的偏移量映射,匹配缓存中的开始时间,计算耗时并记录,最后清理缓存:

import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import java.util.Map;

@Configuration
public class KafkaConfig {

    @Autowired
    private KafkaBatchListener batchListener;

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {

        log.info("创建Kafka监听器工厂");
        ConcurrentKafkaListenerContainerFactory<String, String> factory
                = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setBatchListener(true);
        factory.getContainerProperties().setSyncCommits(false);     
        factory.getContainerProperties()
                .setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        
        // 配置提交回调
        factory.getContainerProperties().setCommitCallback((map, ex) -> {
            long endTime = System.currentTimeMillis();
            if (ex == null) {
                // 遍历提交的每个分区偏移量,计算耗时
                for (Map.Entry<TopicPartition, OffsetAndMetadata> entry : map.entrySet()) {
                    TopicPartition tp = entry.getKey();
                    // 提交的偏移量是下一个要消费的位置,对应已处理的最后一条消息偏移量需减1
                    long committedOffset = entry.getValue().offset() - 1;
                    String key = tp.topic() + "-" + tp.partition() + "-" + committedOffset;
                    Long startTime = batchListener.commitStartTimeCache.remove(key);
                    if (startTime != null) {
                        long cost = endTime - startTime;
                        log.info("分区{}提交耗时:{}ms", tp, cost);
                        // 此处可发布自定义指标,比如Micrometer Timer
                        // timer.record(cost, TimeUnit.MILLISECONDS);
                    }
                }
                log.info("提交成功:{}", map);
            } else {
                // 提交失败时清理对应缓存,避免内存泄漏
                for (Map.Entry<TopicPartition, OffsetAndMetadata> entry : map.entrySet()) {
                    TopicPartition tp = entry.getKey();
                    long committedOffset = entry.getValue().offset() - 1;
                    String key = tp.topic() + "-" + tp.partition() + "-" + committedOffset;
                    batchListener.commitStartTimeCache.remove(key);
                }
                log.error("提交失败:{},异常信息:{}", map, ex.getMessage());
            }
        });
        return factory;
    }

    // 省略consumerFactory()实现
}

关键说明

  • 提交的OffsetAndMetadata中的偏移量是下一个要消费的偏移量,因此已处理的最后一条消息偏移量需要减1,才能和监听器中记录的偏移量匹配。
  • 使用ConcurrentHashMap保证缓存的线程安全,适配异步提交的并发场景。
  • 无论提交成功或失败,都要清理对应缓存,防止内存泄漏。

可选优化

若无需单条消息的耗时统计,仅需批量整体耗时,可简化逻辑:只记录批量的开始时间,用分区+批量最大偏移量作为缓存key即可。

内容的提问来源于stack exchange,提问作者Aakanksha Sharma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 14:57:55