MessageListenerContainer重复指标问题求助:疑为Spring-Kafka问题?
首先得明确:你遇到的重复上报、数值不同的情况,根源是Spring-Kafka的ConcurrentMessageListenerContainer会创建多个独立的KafkaConsumer实例,而不是Kafka本身的问题。让我慢慢给你解释:
为什么会出现多次同类型指标?
当你用ConcurrentMessageListenerContainer时,Spring-Kafka会根据你配置的concurrency参数启动多个consumer线程——你的Topic有3个分区,通常我们会把concurrency设为和分区数一致(也就是3),这样每个consumer对应一个分区。
每个consumer实例都有自己的客户端ID(比如你的输出里的kafka-health-check-0,其实还有-1、-2的实例),而且每个实例都会维护自己独立的metrics统计数据。你看到的三次outgoing_byte_rate,分别是这三个consumer各自的字节发送速率,所以数值不一样,本质上是三个不同实例的指标,不是重复上报。
你的代码逻辑正好验证了这一点
看你贴的consumerMetrics()方法(顺便提一句,代码里有个小语法错误,输出语句末尾多了个+ },应该补上指标值的输出):
public void consumerMetrics() { for (MessageListenerContainer messageListenerContainer : kafkaListenerEndpointRegistry.getListenerContainers()) { Map<String, Map<MetricName, ? extends Metric>> metrics = messageListenerContainer.metrics(); metrics.forEach( (clientid, metricMap) -> { System.out.println("for client id:" + clientid); metricMap.forEach((metricName, metricValue) -> { String cMetricName = metricName.name(); double cMetricValue = metricValue.value(); String cMetricNameClean = cMetricName.replaceAll("[-|.]", "_"); // 修正后的输出语句 System.out.println(clientid + " | " + cMetricNameClean + " : " + cMetricValue); }); }); } }
这段代码会遍历每个listener容器下的所有consumer实例,然后输出每个实例的所有metrics。因为你有3个consumer线程,所以同一种指标会被输出3次,每次对应不同consumer的统计结果。
怎么解决?看你的需求来
1. 我需要区分每个consumer的指标
那其实你现在的输出已经没问题了——每个指标前面带的kafka-health-check-0就是consumer的客户端ID,能清楚区分是哪个实例的指标,只是你误以为是重复的而已。如果觉得混乱,可以在输出里加上线程标识或者更清晰的备注。
2. 我需要整个容器的汇总指标
如果想得到整个listener的总速率(比如所有consumer的outgoing_byte_rate总和),可以对同名称的指标做聚合计算,比如这样改代码:
public void consumerMetrics() { // 用map存聚合后的指标:指标名 -> 总数值 Map<String, Double> aggregatedMetrics = new HashMap<>(); for (MessageListenerContainer container : kafkaListenerEndpointRegistry.getListenerContainers()) { Map<String, Map<MetricName, ? extends Metric>> metrics = container.metrics(); metrics.forEach( (clientid, metricMap) -> { metricMap.forEach((metricName, metricValue) -> { String cleanMetricName = metricName.name().replaceAll("[-|.]", "_"); double metricVal = metricValue.value(); // 累加相同指标的数值 aggregatedMetrics.put(cleanMetricName, aggregatedMetrics.getOrDefault(cleanMetricName, 0.0) + metricVal); }); }); } // 输出汇总后的指标 aggregatedMetrics.forEach((name, totalVal) -> { System.out.println("aggregated | " + name + " : " + totalVal); }); }
这样就能得到所有consumer的同类型指标总和,不会再看到重复的条目。
3. 检查concurrency配置是否合理
如果你不小心把concurrency设成了大于3的值(比如4),那多余的consumer线程会处于空闲状态,但还是会上报自己的metrics。建议确认你的@KafkaListener或者容器配置里的concurrency值,最好和Topic的分区数一致,这样资源利用最合理。
内容的提问来源于stack exchange,提问作者Kenny O'Brien

