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

如何在Kafka ProducerInterceptor中匹配生产记录与对应RecordMetadata

Kafka生产者消息发送&ack时间追踪解决方案

核心适配条件

你选择的自定义ProducerInterceptor方案完全符合需求:无需修改业务代码,仅需要在生产者配置中增加interceptor.classes参数指定自定义拦截器全类名即可生效。

方案1:基于分区有序性维护映射队列(生产环境首选)

Kafka对单分区消息有严格的顺序保证:同一个分区下的消息发送顺序和ack回调顺序完全一致,基于这个特性可以实现消息的精准匹配,实现逻辑如下:

  1. 自定义拦截器内维护线程安全的存储结构:ConcurrentHashMap<String, ConcurrentLinkedQueue<Pair<String, Long>>> ,key为topic:partition格式的分区唯一标识,value为有序队列,队列存储元素为<消息唯一ID, 发送时间戳>的键值对
  2. onSend方法逻辑:
    • 从ProducerRecord的header中取出预设的唯一字符串ID
    • 记录当前系统时间作为消息发送时间
    • 拿到消息分配的分区号,构造topic:partitionkey,将<消息ID, 发送时间>追加到对应队列的尾部
  3. 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扩展+异步消费匹配(无顺序依赖)

如果你的场景下分区顺序可能被打乱(如未开幂等且允许重试乱序),可以用该方案做兜底:

  1. 在onSend方法中,除了原有唯一ID外,额外往ProducerRecord的header中塞入当前发送时间戳
  2. 单独启动一个异步消费线程,批量拉取对应topic的消息,从header中取出唯一ID和发送时间,结合Broker端的写入时间(将topic配置message.timestamp.type设置为LogAppendTime后,消息自带的timestamp即为Broker存储时间)做匹配

该方案完全不依赖消息顺序,但是会有一定的延迟,且需要额外的消费资源,适合对实时性要求不高的场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 18:03:05