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

Spring Cloud Stream Kafka生产者线程持续增长问题

Spring Cloud Stream Kafka Binder 同应用多消费者场景下生产者线程泄漏问题

问题表现

应用基于Spring Cloud Stream Kafka Binder实现多级消息处理链路:

  • 消费源主题消息,处理后输出到多个中间主题
  • 中间主题由同一应用实例消费,处理完成后最终输出到结果主题

观测到的异常特征:

  • 每当首个消费者消费新消息时,系统会创建大量Kafka生产者相关线程
  • 这些线程创建后持续存活不会被回收,随消息消费次数增加持续累积
  • 复现边界清晰:将首个消费者schedulingConsumer与其余两个消费者consumerSearch1、consumerSearch2拆分独立部署后问题不再复现,仅当所有消费者运行在同一应用实例中时异常触发
  • 线程栈可观测到大量命名格式为kafka-producer-network-thread | producer-*的常驻线程

对应配置

spring:
  cloud:
    stream:
      function:
        definition: schedulingConsumer;consumerSearch1;consumerSearch2
      default:
        group: ${kafka.group}
        contentType: application/json
        consumer:
          maxAttempts: 1
          backOffMaxInterval: 30
          retryableExceptions:
            org.springframework.messaging.converter.MessageConversionException: false
      kafka:
        binder:
          brokers: ${kafka.brokers}
          headerMapperBeanName: kafkaHeaderMapper
          producerProperties:
            linger.ms: 500
            batch.size: ${kafka.batchs.size}
            compression.type: gzip
          consumerProperties:
            session.timeout.ms: ${kafka.session.timeout.ms}
            max.poll.interval.ms: ${kafka.poll.interval}
            max.poll.records: ${kafka.poll.records}
            commit.interval.ms: 500
            allow.auto.create.topics: false
      bindings:
        schedulingConsumer-in-0:
          destination: ${kafka.topics.schedules}
          consumer.concurrency: 5

        search1-out:
          destination: ${kafka.topics.groups.search1}
        search2-out:
          destination: ${kafka.topics.groups.search2}

        consumerSearch1-in-0:
          destination: ${kafka.topics.groups.search1}
        consumerSearch2-in-0:
          destination: ${kafka.topics.groups.search2}

        datasource-out:
          destination: ${kafka.topics.search.output}

根因说明

该问题是Spring Cloud Stream 3.2.4之前版本Kafka Binder的已知缺陷:
当同一应用上下文中存在多个函数绑定、且入口消费者配置了大于1的消费并发(当前配置中schedulingConsumer并发为5)时,binder内部生产者缓存的key计算逻辑存在bug,无法命中已创建的可复用生产者实例,会为每个消费线程、每个输出目标重复创建独立的KafkaProducer实例。
每个KafkaProducer实例初始化时会创建独立的网络IO线程、批处理发送线程,这些重复创建的生产者实例既没有被缓存复用,也没有被Spring上下文纳入生命周期管理,不会在消息处理完成后调用close()方法回收资源,最终导致相关线程持续累积无法回收。

解决方案

可根据实际场景选择以下任意一种方案修复:

  • 版本升级(推荐):将Spring Cloud Stream版本升级到3.2.4及以上,官方已修复多函数场景下生产者缓存key计算错误的问题,相同配置的生产者实例会被正确复用,不会重复创建。
  • 配置规避(无法升级版本时使用):在binder配置中显式开启生产者全局复用,添加如下配置:
    spring:
      cloud:
        stream:
          kafka:
            binder:
              # 关闭按绑定独立创建生产者的默认逻辑
              producerPerBinding: false
    
    注意如果代码中使用StreamBridge发送消息,需注入Spring上下文默认的单例ProducerFactory,不要手动实例化新的生产者工厂。
  • 架构拆分:即已经验证过的部署方式,将不同消费逻辑的函数拆分到独立的应用进程部署,从物理上隔离绑定上下文,避免缓存逻辑冲突。

内容的提问来源于stack exchange,提问作者Jadest

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 10:45:38