如何聚合Kafka Streams应用中Micrometer计数器数组的实时指标值
问题根因
原有配置Bean不生效的核心原因是Bean初始化阶段Kafka Streams线程还未完成启动,对应的线程维度指标还未注册到MeterRegistry,此时采集到的FunctionCounter列表为空,且后续新增的线程指标也不会自动加入初始化时绑定的state对象,因此无法得到正确的聚合结果。
修正实现方案
将聚合逻辑放到FunctionCounter的取值函数中,每次指标被查询时动态从注册中心拉取符合条件的指标实时计算,天然适配底层指标更新、Kafka Streams线程动态增减的场景,无需额外实现更新触发机制。
@Configuration @Slf4j public class MetricsConfig { @Bean public FunctionCounter getAggregateCounter(MeterRegistry registry) { return FunctionCounter .builder("Combined_Output_Message_Count", registry, reg -> reg.getMeters().stream() .filter(meter -> meter.getId().getName().startsWith("Output_Message_Count")) .filter(FunctionCounter.class::isInstance) .map(FunctionCounter.class::cast) .mapToDouble(FunctionCounter::count) .sum() ) .description("Kafka Streams多线程输出消息数聚合总和") .tags("region", "test") .register(registry); } }
如果需要聚合kafka_stream_thread_task_created_total这类指标,只需将startsWith的参数替换为对应指标名即可,最终聚合结果符合预期(示例数据聚合后数值为20)。
内容的提问来源于stack exchange,提问作者K Olusanya
相关产品推荐
相关产品推荐

