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,两种修复方式:- 升级Spring Cloud Stream到3.2.x及以上版本,官方已经修复了绑定失败时缓存无效通道的问题
- 不方便升级的话,可以在应用启动后主动调用一次StreamBridge发送测试消息,绑定成功后再对外提供服务,避免运行时首次请求触发异常
内容的提问来源于stack exchange,提问作者Bert S.
相关产品推荐
相关产品推荐

