Spring调度器+Kafka集成时Brave Baggage引发OOM问题求助
问题描述
应用使用Kafka Streams及Spring标准调度器,调度器运行时会向Kafka发送并消费消息。堆转储分析显示Brave Baggage数据被大量复制,关闭调度器后堆内存停止增长。此外日志中持续出现Kafka元数据重置信息,关闭Prometheus后该日志消失。
相关配置
kafka: streams.binder.serdeError: sendToDlq binder: brokers: localhost:9092 auto-create-topics: true auto-add-partitions: true admin: fail-fast: true bootstrap-servers: localhost:9092 cloud: stream: bindings: merchant: destination: ******* legal: destination: ******* account: destination: ******* merchant_session: destination: ******* #async-task in-queue: destination: ********* contentType: application/json group: merchant consumer: max-attempts: 1 partitioned: true concurrency: 10 in-queue-dlq: destination: ********** contentType: application/json group: merchant consumer: max-attempts: 1 back-off-initial-interval: 5000 out-queue: destination: ********** contentType: application/json out-queue-dlq: destination: ********** contentType: application/json kafka: binder: auto-add-partitions: true auto-create-topics: true
调度器代码
@Transactional @Scheduled( cron = "${task-processor.cronExpression:*/5 * * * * ?}" ) public void rescheduleTask() { Integer taskRetrySize = this.properties.getTaskRetrySize(); List list = this.databaseDataProvider.findPlannedTasks(taskRetrySize); if (!list.isEmpty()) { log.debug("Rescheduling planned tasks. Found {} tasks for rescheduling.", list.size()); } IMessageQueueDataProvider var10001 = this.messageQueueDataProvider; Objects.requireNonNull(var10001); list.forEach(var10001::resendRequest); }
异常日志
2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-0 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-5 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-8 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-2 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-9 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-1 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-4 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-6 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-7 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.090 [pool-21-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.execute-3 to 0 since the associated topicId changed from null to X41sXMr6TqStCZ5BV-8Lhg 2023-02-27 22:13:25.096 [pool-22-thread-1] [- | -] INFO org.apache.kafka.clients.Metadata Metadata.updateLatestMetadata(...) (Metadata.java:401) - [Consumer clientId=consumer-merchant-14, groupId=merchant] Resetting the last seen epoch of partition rdbo.merchant.task.error-0 to 0 since the associated topicId changed from null to _E_0g0iZRGeBSq0PWrpGtQ 2023-02-27 22:13:25.101 [pool-22-thread-1] [- | -] INFO o.a.k.c.c.i.ConsumerCoordinator
原因分析与修复方案
1. Brave Baggage内存增长问题
- 核心原因:Spring调度器线程未正确清理Brave的Baggage上下文。每次调度任务执行时,Brave会为当前线程创建新的Baggage实例,但任务结束后上下文未被清除,导致实例堆积;同时Kafka消息发送/消费时会传播Baggage,进一步加剧复制。
- 修复步骤:
- 在调度任务方法末尾手动清理Brave上下文:
@Transactional @Scheduled(cron = "${task-processor.cronExpression:*/5 * * * * ?}") public void rescheduleTask() { try { Integer taskRetrySize = this.properties.getTaskRetrySize(); List list = this.databaseDataProvider.findPlannedTasks(taskRetrySize); if (!list.isEmpty()) { log.debug("Rescheduling planned tasks. Found {} tasks for rescheduling.", list.size()); } IMessageQueueDataProvider messageQueueDataProvider = this.messageQueueDataProvider; Objects.requireNonNull(messageQueueDataProvider); list.forEach(messageQueueDataProvider::resendRequest); } finally { // 清理Brave上下文 BaggageContext.clear(); CurrentTraceContext.clear(); } } - 检查Brave配置,禁用不必要的Baggage字段,减少复制的数据量。
- 在调度任务方法末尾手动清理Brave上下文:
2. Kafka元数据重置日志问题
- 核心原因:Prometheus监控可能会频繁调用Kafka的AdminClient获取元数据,导致Kafka客户端反复更新topicId(尤其是自动创建的topic),触发元数据重置日志。
- 修复步骤:
- 调整Prometheus监控策略,避免频繁拉取Kafka元数据,或禁用对Kafka AdminClient的监控。
- 提前手动创建所需的Kafka topic,关闭自动创建topic的配置,避免动态创建时的元数据波动:
kafka: binder: auto-create-topics: false cloud: stream: kafka: binder: auto-create-topics: false - 升级Kafka客户端版本,部分旧版本存在元数据更新的bug,新版本已修复。
内容的提问来源于stack exchange,提问作者Олег Головин
相关产品推荐
相关产品推荐

