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

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字段,减少复制的数据量。

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,提问作者Олег Головин

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 09:45:07