如何在Kafka ProducerInterceptor中匹配生产记录与对应RecordMetadata
Kafka生产者消息发送&ack时间追踪解决方案
核心适配条件
你选择的自定义ProducerInterceptor方案完全符合需求:无需修改业务代码,仅需要在生产者配置中增加interceptor.classes参数指定自定义拦截器全类名即可生效。
方案1:基于分区有序性维护映射队列(生产环境首选)
Kafka对单分区消息有严格的顺序保证:同一个分区下的消息发送顺序和ack回调顺序完全一致,基于这个特性可以实现消息的精准匹配,实现逻辑如下:
- 自定义拦截器内维护线程安全的存储结构:
ConcurrentHashMap<String, ConcurrentLinkedQueue<Pair<String, Long>>>,key为topic:partition格式的分区唯一标识,value为有序队列,队列存储元素为<消息唯一ID, 发送时间戳>的键值对 onSend方法逻辑:- 从ProducerRecord的header中取出预设的唯一字符串ID
- 记录当前系统时间作为消息发送时间
- 拿到消息分配的分区号,构造
topic:partitionkey,将<消息ID, 发送时间>追加到对应队列的尾部
onAcknowledgement方法逻辑:- 从RecordMetadata中获取topic和分区号,构造和onSend中一致的key
- 从对应队列头部取出第一条元素,即为当前ack消息对应的ID和发送时间
- 记录当前系统时间(或使用RecordMetadata的timestamp,按需选择是客户端收到ack的时间还是Broker端写入时间)作为ack时间
- 完成匹配后即可做追踪数据的上报、存储等后续逻辑
- 如果exception不为空,说明消息发送失败,同样取出队列头元素做异常打点即可
注意事项:
- 如果生产者开启了重试,建议同时配置
enable.idempotence=true开启幂等性,保证重试场景下消息顺序不会乱,避免队列匹配错位- 建议给队列中的元素增加过期时间,定期清理超过阈值(如10分钟)还未被消费的元素,避免极端情况下消息丢失导致的内存泄漏
代码示例
import org.apache.kafka.clients.producer.*; import org.apache.kafka.common.header.Header; import org.apache.commons.lang3.tuple.ImmutablePair; import java.nio.charset.StandardCharsets; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentLinkedQueue; public class TracingProducerInterceptor<K, V> implements ProducerInterceptor<K, V> { private final ConcurrentHashMap<String, ConcurrentLinkedQueue<ImmutablePair<String, Long>>> partitionQueueMap = new ConcurrentHashMap<>(); // 替换为你实际存消息唯一ID的header key private static final String MSG_UNIQUE_ID_HEADER = "msg_unique_id"; @Override public ProducerRecord<K, V> onSend(ProducerRecord<K, V> record) { String msgId = null; for (Header header : record.headers()) { if (MSG_UNIQUE_ID_HEADER.equals(header.key())) { msgId = new String(header.value(), StandardCharsets.UTF_8); break; } } if (msgId != null) { String partitionKey = record.topic() + ":" + record.partition(); partitionQueueMap.computeIfAbsent(partitionKey, k -> new ConcurrentLinkedQueue<>()) .add(ImmutablePair.of(msgId, System.currentTimeMillis())); } return record; } @Override public void onAcknowledgement(RecordMetadata metadata, Exception exception) { String partitionKey = metadata.topic() + ":" + metadata.partition(); ConcurrentLinkedQueue<ImmutablePair<String, Long>> queue = partitionQueueMap.get(partitionKey); if (queue == null || queue.isEmpty()) { return; } ImmutablePair<String, Long> msgInfo = queue.poll(); if (msgInfo == null) { return; } String msgId = msgInfo.left; Long sendTime = msgInfo.right; long ackTime = System.currentTimeMillis(); // 此处添加你的追踪逻辑:上报、存储均可 if (exception != null) { // 消息发送失败的处理逻辑 } } @Override public void close() { partitionQueueMap.clear(); } @Override public void configure(Map<String, ?> configs) { // 自定义配置初始化逻辑 } }
方案2:Header扩展+异步消费匹配(无顺序依赖)
如果你的场景下分区顺序可能被打乱(如未开幂等且允许重试乱序),可以用该方案做兜底:
- 在
onSend方法中,除了原有唯一ID外,额外往ProducerRecord的header中塞入当前发送时间戳 - 单独启动一个异步消费线程,批量拉取对应topic的消息,从header中取出唯一ID和发送时间,结合Broker端的写入时间(将topic配置
message.timestamp.type设置为LogAppendTime后,消息自带的timestamp即为Broker存储时间)做匹配
该方案完全不依赖消息顺序,但是会有一定的延迟,且需要额外的消费资源,适合对实时性要求不高的场景
内容的提问来源于stack exchange,提问作者dmorar
相关产品推荐
相关产品推荐

