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
相关产品推荐
相关产品推荐

