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

如何结合Micrometer Timer统计KafkaListener消费消息内元素的数量与耗时

方案1:业务代码手动埋点(单监听器场景推荐,实现简单灵活)

直接在消费逻辑中引入Spring Boot自带的Micrometer MeterRegistry 自定义你需要的统计维度,代码示例如下:

import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.Timer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import java.util.List;

@Component
public class DataConsumer {

    @Autowired
    private MeterRegistry meterRegistry;

    @KafkaListener(
        id = "dataConsumer",
        topics = "data.topic",
        groupId = "${spring.kafka.consumer.group-id}",
        containerFactory = "dataKafkaListenerContainerFactory")
    public void consumeData(DataContainer message) {
        List<Data> dataList = message.getList();
        int elementCount = dataList.size();
        // 统计处理的元素总数量,计数维度为单条数据元素
        meterRegistry.counter("kafka.listener.data.element.total", 
                "listenerId", "dataConsumer",
                "topic", "data.topic")
                .increment(elementCount);
        
        // 启动计时,统计本次消息处理总耗时
        Timer.Sample processTimer = Timer.start(meterRegistry);
        try {
            // 原有业务处理逻辑
            doProcess(dataList);
        } finally {
            // 结束计时,记录总耗时
            processTimer.stop(Timer.builder("kafka.listener.data.process.duration")
                    .tag("listenerId", "dataConsumer")
                    .tag("topic", "data.topic")
                    .register(meterRegistry));
            
            // 可选:统计单元素平均处理耗时,单位为毫秒
            if (elementCount > 0) {
                long totalNanos = processTimer.stop(Timer.builder("tmp").register(meterRegistry));
                double perElementMs = totalNanos / 1_000_000.0 / elementCount;
                meterRegistry.summary("kafka.listener.data.per.element.duration.ms",
                        "listenerId", "dataConsumer")
                        .record(perElementMs);
            }
        }
    }

    private void doProcess(List<Data> dataList) {
        // 原有业务逻辑
    }
}

方案2:全局拦截器统一埋点(多监听器场景推荐,统一规范)

如果有多个同类型的列表消息监听器,可自定义RecordInterceptor实现全局埋点,避免重复代码:

  1. 实现RecordInterceptor接口,在onMessage方法前后插入统计逻辑,从反序列化后的消息体中获取列表长度记录指标
  2. 将自定义拦截器配置到dataKafkaListenerContainerFactory的recordInterceptor属性中即可全局生效

指标查询说明

  • 累计处理元素总数量:访问/actuator/metrics/kafka.listener.data.element.total?tag=listenerId:dataConsumer,返回的COUNT值即为所有消息内元素的总数
  • 累计处理总耗时:访问/actuator/metrics/kafka.listener.data.process.duration?tag=listenerId:dataConsumer,返回的TOTAL值即为所有消息的处理总耗时
  • 你也可以根据自己的需求添加groupId、实例ip等自定义tag,方便多部署环境下的指标聚合

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 08:57:00