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

Spring Cloud Stream Kafka事务生产者发送高分区主题超时如何解决

问题描述

我使用Spring Cloud Stream Kafka binder配置了事务生产者,通过StreamBridge生产消息,其中有一个主题的分区数高达135,使用过程中遇到如下问题:系统仅在首次向流发送消息时才会尝试创建事务生产者,此时会尝试获取每个分区的信息但操作超时,之后系统会缓存一个fullChannelName为unknown.channel.name的MessageChannel,导致该通道后续所有StreamBridge.send调用都失败。请问有什么方案可以正常发送消息?是否可以在应用启动时就预先拉取分区数据,避免首次发送时的超时问题?

报错日志

1. 每个分区重复触发的报错日志

2021-08-24 16:45:01.160 ERROR 59629 --- [   scheduling-1] o.s.c.s.b.k.p.KafkaTopicProvisioner      : Failed to obtain partition information

org.springframework.kafka.KafkaException: initTransactions() failed; nested exception is org.apache.kafka.common.errors.TimeoutException: Timeout expired after 60000 milliseconds while awaiting InitProducerId
    at org.springframework.kafka.core.DefaultKafkaProducerFactory.doCreateTxProducer(DefaultKafkaProducerFactory.java:732) ~[spring-kafka-2.7.6.jar:2.7.6]
    at org.springframework.kafka.core.DefaultKafkaProducerFactory.createTransactionalProducer(DefaultKafkaProducerFactory.java:673) ~[spring-kafka-2.7.6.jar:2.7.6]
    at org.springframework.kafka.core.DefaultKafkaProducerFactory.createTransactionalProducerForPartition(DefaultKafkaProducerFactory.java:594) ~[spring-kafka-2.7.6.jar:2.7.6]
    at org.springframework.kafka.core.DefaultKafkaProducerFactory.doCreateProducer(DefaultKafkaProducerFactory.java:530) ~[spring-kafka-2.7.6.jar:2.7.6]
    at org.springframework.kafka.core.DefaultKafkaProducerFactory.createProducer(DefaultKafkaProducerFactory.java:519) ~[spring-kafka-2.7.6.jar:2.7.6]
    at org.springframework.kafka.core.DefaultKafkaProducerFactory.createProducer(DefaultKafkaProducerFactory.java:513) ~[spring-kafka-2.7.6.jar:2.7.6]
    at org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.lambda$createProducerMessageHandler$0(KafkaMessageChannelBinder.java:396) ~[spring-cloud-stream-binder-kafka-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner.lambda$getPartitionsForTopic$6(KafkaTopicProvisioner.java:535) ~[spring-cloud-stream-binder-kafka-core-3.1.3.jar:3.1.3]

2. Binder相关报错日志

2021-08-24 16:45:12.768 ERROR 59629 --- [   scheduling-1] o.s.cloud.stream.binding.BindingService  : Failed to create producer binding; retrying in 30 seconds

org.springframework.cloud.stream.binder.BinderException: Cannot initialize binder:
    at org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner.getPartitionsForTopic(KafkaTopicProvisioner.java:591) ~[spring-cloud-stream-binder-kafka-core-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.createProducerMessageHandler(KafkaMessageChannelBinder.java:394) ~[spring-cloud-stream-binder-kafka-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.createProducerMessageHandler(KafkaMessageChannelBinder.java:158) ~[spring-cloud-stream-binder-kafka-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.AbstractMessageChannelBinder.doBindProducer(AbstractMessageChannelBinder.java:226) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.AbstractMessageChannelBinder.doBindProducer(AbstractMessageChannelBinder.java:91) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.AbstractBinder.bindProducer(AbstractBinder.java:152) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binding.BindingService.lambda$rescheduleProducerBinding$4(BindingService.java:343) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.scheduling.support.DelegatingErrorHandlingRunnable.run(DelegatingErrorHandlingRunnable.java:54) ~[spring-context-5.3.9.jar:5.3.9]
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) ~[na:na]
    at java.base/java.util.concurrent.FutureTask.run$$$capture(FutureTask.java:264) ~[na:na]
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java) ~[na:na]
    at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:829) ~[na:na]

3. 反复出现的初始化错误日志

2021-08-24 17:08:02.753 ERROR 59629 --- [   scheduling-1] o.s.c.s.b.k.p.KafkaTopicProvisioner      : Cannot initialize Binder

java.lang.IllegalStateException: The number of expected partitions was: 1, but 0 has been found instead
    at org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner.lambda$getPartitionsForTopic$6(KafkaTopicProvisioner.java:579) ~[spring-cloud-stream-binder-kafka-core-3.1.3.jar:3.1.3]
    at org.springframework.retry.support.RetryTemplate.doExecute(RetryTemplate.java:329) ~[spring-retry-1.3.1.jar:na]
    at org.springframework.retry.support.RetryTemplate.execute(RetryTemplate.java:209) ~[spring-retry-1.3.1.jar:na]
    at org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner.getPartitionsForTopic(KafkaTopicProvisioner.java:530) ~[spring-cloud-stream-binder-kafka-core-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.createProducerMessageHandler(KafkaMessageChannelBinder.java:394) ~[spring-cloud-stream-binder-kafka-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.createProducerMessageHandler(KafkaMessageChannelBinder.java:158) ~[spring-cloud-stream-binder-kafka-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.AbstractMessageChannelBinder.doBindProducer(AbstractMessageChannelBinder.java:226) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.AbstractMessageChannelBinder.doBindProducer(AbstractMessageChannelBinder.java:91) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binder.AbstractBinder.bindProducer(AbstractBinder.java:152) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.binding.BindingService.lambda$rescheduleProducerBinding$4(BindingService.java:343) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.scheduling.support.DelegatingErrorHandlingRunnable.run(DelegatingErrorHandlingRunnable.java:54) ~[spring-context-5.3.9.jar:5.3.9]
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) ~[na:na]
    at java.base/java.util.concurrent.FutureTask.run$$$capture(FutureTask.java:264) ~[na:na]
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java) ~[na:na]
    at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:829) ~[na:na]

4. 最终消息发送失败报错日志

2021-08-24 16:59:43.224 ERROR 59629 --- [nio-8080-exec-4] s.s.e.RestResponseEntityExceptionHandler : org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'unknown.channel.name'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[235], headers={<snip>}, ce-source=https://spring.io/, kafka_messageKey=d15140aa-3b64-416a-9b1f-8829ee7c8319, ce-type=<my type>, message-type=cloudevent, specversion=1.0.8, id=c05e84b6-40ad-3787-32b8-54096cb735aa, contentType=application/avro, timestamp=1629841968540}], failedMessage=GenericMessage [payload=byte[235], headers={<snip>, contentType=application/avro, timestamp=1629841968540}]

org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'unknown.channel.name'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[235], headers={<snip>, contentType=application/avro, timestamp=1629841968540}]
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:76) ~[spring-integration-core-5.5.3.jar:5.5.3]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:317) ~[spring-integration-core-5.5.3.jar:5.5.3]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272) ~[spring-integration-core-5.5.3.jar:5.5.3]
    at org.springframework.cloud.stream.function.StreamBridge.send(StreamBridge.java:215) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
    at org.springframework.cloud.stream.function.StreamBridge.send(StreamBridge.java:156) ~[spring-cloud-stream-3.1.3.jar:3.1.3]
...my app stacktrace...
解决方案
  • 预初始化通道,避免懒加载超时
    可以在应用启动时提前完成生产者绑定,不需要等到首次发送消息时才执行初始化。只需要在配置文件中显式声明输出绑定,指定对应主题的分区数即可:
spring:
  cloud:
    stream:
      bindings:
        # 替换为实际使用的绑定名,命名规则为<自定义名>-out-<序号>
        custom-output-out-0:
          destination: 实际Kafka主题名
          producer:
            # 直接填主题实际分区数135,避免运行时拉取分区信息
            partition-count: 135

配置后应用启动时就会主动完成通道绑定、元数据拉取、事务生产者初始化,完全避免首次发送时的超时问题。

  • 优化事务生产者配置,降低初始化耗时
    本次超时的核心原因是默认事务生产者按分区隔离,135个分区会触发135次InitProducerId请求到Kafka集群,很容易触发超时。可以关闭分区隔离特性,所有分区复用同一个事务生产者:
spring:
  cloud:
    stream:
      kafka:
        bindings:
          custom-output-out-0:
            producer:
              transactional:
                # 关闭按分区隔离事务生产者
                per-partition-bound: false
        binder:
          operation-timeout: 120000
          configuration:
            # 调大生产者阻塞操作超时时间
            max.block.ms: 120000

调整后仅需要1次InitProducerId请求,初始化耗时会大幅降低。

  • 修复无效通道缓存问题
    你遇到的unknown.channel.name缓存是Spring Cloud Stream 3.1.x版本的已知bug,两种修复方式:
    1. 升级Spring Cloud Stream到3.2.x及以上版本,官方已经修复了绑定失败时缓存无效通道的问题
    2. 不方便升级的话,可以在应用启动后主动调用一次StreamBridge发送测试消息,绑定成功后再对外提供服务,避免运行时首次请求触发异常

内容的提问来源于stack exchange,提问作者Bert S.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 04:15:02