如何减少Spring-Kafka消费者的micrometer-kafka-metrics线程数量?
我们基于Spring Boot开发的应用,使用Spring Kafka的@KafkaListener实现了同消费组下的多消费者。排查发现每个消费者会创建4类线程,本次聚焦micrometer-kafka-metrics线程:
该线程是由io.micrometer:micrometer.core在KafkaMetrics中创建的守护线程,核心代码如下:
ScheduledExecutorService scheduler = Executors .newSingleThreadScheduledExecutor(new NamedThreadFactory("micrometer-kafka-metrics"));
调用链路为:spring-boot-actuator-autoconfigure中的KafkaMetricsAutoConfiguration会为应用创建单个MicrometerConsumerListener实例,每当新增Consumer时,该实例的consumerAdded(..)方法会被调用,创建继承自KafkaMetrics的KafkaClientMetrics,最终导致micrometer-kafka-metrics线程数与消费者数量完全一致。
若消费者工厂并发数为N(通过new ConcurrentKafkaListenerContainerFactory<>().setConcurrency(N);配置),@KafkaListener数量为M(对应M个主题),则会生成N*M个这类线程,部分场景下可达数百个。但这类线程本身轻量,默认1分钟执行一次,因此希望在保留原有并发消费者数量的前提下,减少Metrics线程的数量。
理想方案是支持传入可共享的自定义线程池,但目前Micrometer似乎无此能力,是否有其他可行方案?
方案1:自定义MicrometerConsumerListener替换默认实现
直接替换Spring Boot自动配置的MicrometerConsumerListener,实现共享线程池的逻辑:
- 先禁用自动配置的
KafkaMetricsAutoConfiguration,在启动类添加:
@SpringBootApplication(exclude = KafkaMetricsAutoConfiguration.class)
- 自定义
MicrometerConsumerListener,使用共享线程池:
@Component public class SharedPoolMicrometerConsumerListener implements ConsumerListener { private final MeterRegistry meterRegistry; private final ScheduledExecutorService sharedScheduler; public SharedPoolMicrometerConsumerListener(MeterRegistry meterRegistry) { this.meterRegistry = meterRegistry; // 按需配置共享线程池大小,比如2-4个核心线程 this.sharedScheduler = Executors.newScheduledThreadPool(2, new NamedThreadFactory("shared-micrometer-kafka-metrics")); } @Override public void consumerAdded(String id, Consumer<?, ?> consumer) { // 重写KafkaClientMetrics的线程池创建逻辑,传入共享池 new KafkaClientMetrics(consumer, meterRegistry) { @Override protected ScheduledExecutorService createScheduler() { return sharedScheduler; } }.bindTo(meterRegistry); } @Override public void consumerRemoved(String id, Consumer<?, ?> consumer) { KafkaClientMetrics.unregisterMetrics(consumer, meterRegistry); } @PreDestroy public void shutdown() { sharedScheduler.shutdown(); } }
所有消费者会复用同一个共享线程池,线程数不再随消费者数量增长,由你配置的线程池大小决定。
方案2:反射替换KafkaMetrics的线程池(侵入性方案)
如果不想禁用自动配置,可以通过反射替换KafkaMetrics内部的线程池,注意该方案依赖Micrometer内部实现,升级版本时需验证兼容性:
@Component public class KafkaMetricsThreadPoolCustomizer implements ApplicationListener<ApplicationReadyEvent> { private final ScheduledExecutorService sharedScheduler = Executors.newScheduledThreadPool(3, new NamedThreadFactory("shared-micrometer-kafka-metrics")); @Override public void onApplicationEvent(ApplicationReadyEvent event) { // 遍历所有KafkaClientMetrics实例,替换线程池 Metrics.globalRegistry.getMeters().stream() .filter(meter -> meter instanceof KafkaClientMetrics) .map(meter -> (KafkaClientMetrics) meter) .forEach(metrics -> { try { Field schedulerField = KafkaMetrics.class.getDeclaredField("scheduler"); schedulerField.setAccessible(true); ScheduledExecutorService oldScheduler = (ScheduledExecutorService) schedulerField.get(metrics); oldScheduler.shutdown(); schedulerField.set(metrics, sharedScheduler); } catch (NoSuchFieldException | IllegalAccessException e) { // 记录异常日志 e.printStackTrace(); } }); } @PreDestroy public void shutdown() { sharedScheduler.shutdown(); } }
方案3:关闭Kafka消费者的Micrometer Metrics(非必需场景)
如果这类Metrics不是业务必需的,可以直接关闭:
在application.yml中添加配置:
management: metrics: enable: kafka: false
或者直接禁用KafkaMetricsAutoConfiguration,这样不会创建任何micrometer-kafka-metrics线程,但会丢失Kafka消费者的Metrics数据,需根据实际需求判断。
内容的提问来源于stack exchange,提问作者Vasyl Sarzhynskyi

