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: falseStreamBridge发送消息,需注入Spring上下文默认的单例ProducerFactory,不要手动实例化新的生产者工厂。 - 架构拆分:即已经验证过的部署方式,将不同消费逻辑的函数拆分到独立的应用进程部署,从物理上隔离绑定上下文,避免缓存逻辑冲突。
内容的提问来源于stack exchange,提问作者Jadest
相关产品推荐
相关产品推荐

