SCDF流数据管道突发中断求助(关联Kafka及Spring异常)
问题描述
我们运行一个由Kafka主题(AWS-MSK)触发的流数据管道,原本运行正常,却突然中断并抛出如下异常。销毁重建管道后能正常运行,但一段时间后会再次中断。
org.springframework.kafka.listener.ListenerExecutionFailedException: Listener failed; nested exception is org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'application.transactionProcessor-in-0'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[77], headers={kafka_offset=191, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@3ea3db9e, deliveryAttempt=3, kafka_timestampType=CREATE_TIME, kafka_receivedPartitionId=0, kafka_receivedTopic=aws-db4.BW3.INT_TRANS, kafka_receivedTimestamp=1661930907011, contentType=application/json, kafka_groupId=anonymous.6d66497c-d3fe-4109-9a7a-39b5aaddf974}], failedMessage=GenericMessage [payload=byte[77], headers={kafka_offset=191, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@3ea3db9e, deliveryAttempt=3, kafka_timestampType=CREATE_TIME, kafka_receivedPartitionId=0, kafka_receivedTopic=aws-db4.BW3.INT_TRANS, kafka_receivedTimestamp=1661930907011, contentType=application/json, kafka_groupId=anonymous.6d66497c-d3fe-4109-9a7a-39b5aaddf974}] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.decorateException(KafkaMessageListenerContainer.java:2683) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2649) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeOnMessage(KafkaMessageListenerContainer.java:2609) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:2536) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeWithRecords(KafkaMessageListenerContainer.java:2427) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:2305) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:1979) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1364) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1355) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1247) ~[spring-kafka-2.8.4.jar:2.8.4] at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source) ~[na:na] at java.base/java.util.concurrent.FutureTask.run(Unknown Source) ~[na:na] at java.base/java.lang.Thread.run(Unknown Source) ~[na:na] Caused by: org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'application.transactionProcessor-in-0'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[77], headers={kafka_offset=191, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@3ea3db9e, deliveryAttempt=3, kafka_timestampType=CREATE_TIME, kafka_receivedPartitionId=0, kafka_receivedTopic=aws-db4.BW3.INT_TRANS, kafka_receivedTimestamp=1661930907011, contentType=application/json, kafka_groupId=anonymous.6d66497c-d3fe-4109-9a7a-39b5aaddf974}] at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:76) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:317) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:272) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187) ~[spring-messaging-5.3.18.jar:5.3.18] at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166) ~[spring-messaging-5.3.18.jar:5.3.18] at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) ~[spring-messaging-5.3.18.jar:5.3.18] at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109) ~[spring-messaging-5.3.18.jar:5.3.18] at org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:216) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.sendMessageIfAny(KafkaMessageDrivenChannelAdapter.java:397) ~[spring-integration-kafka-5.5.10.jar:5.5.10] at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.access$300(KafkaMessageDrivenChannelAdapter.java:83) ~[spring-integration-kafka-5.5.10.jar:5.5.10] at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationRecordMessageListener.onMessage(KafkaMessageDrivenChannelAdapter.java:454) ~[spring-integration-kafka-5.5.10.jar:5.5.10] at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationRecordMessageListener.onMessage(KafkaMessageDrivenChannelAdapter.java:428) ~[spring-integration-kafka-5.5.10.jar:5.5.10] at org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter.lambda$onMessage$0(RetryingMessageListenerAdapter.java:125) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.retry.support.RetryTemplate.doExecute(RetryTemplate.java:329) ~[spring-retry-1.3.2.jar:na] at org.springframework.retry.support.RetryTemplate.execute(RetryTemplate.java:255) ~[spring-retry-1.3.2.jar:na] at org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter.onMessage(RetryingMessageListenerAdapter.java:119) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter.onMessage(RetryingMessageListenerAdapter.java:42) ~[spring-kafka-2.8.4.jar:2.8.4] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2629) ~[spring-kafka-2.8.4.jar:2.8.4] ... 11 common frames omitted Caused by: org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:139) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:72) ~[spring-integration-core-5.5.10.jar:5.5.10] ... 28 common frames omitted
问题分析与解决方案
核心问题定位
异常的核心是Dispatcher has no subscribers for channel 'application.transactionProcessor-in-0',说明Kafka消费的消息被发送到application.transactionProcessor-in-0通道后,没有任何订阅者接收处理,最终导致管道中断。
可能原因及解决办法
订阅者Bean意外销毁
检查Spring容器中处理该通道消息的Bean(如标注@Service、@Component的处理器类)是否被意外销毁:- 排查Bean的作用域配置,避免使用
prototype作用域导致Bean频繁创建销毁 - 检查自定义的
BeanPostProcessor或生命周期回调逻辑,确保没有错误触发Bean销毁 - 监控JVM内存使用情况,排查是否因为内存不足导致GC误回收Bean(可调整JVM参数优化内存管理)
- 排查Bean的作用域配置,避免使用
通道与订阅者绑定失效
确认通道与订阅者的绑定关系是否正确:- 检查处理器方法是否通过
@ServiceActivator(inputChannel = "application.transactionProcessor-in-0")正确关联目标通道 - 如果使用Java DSL配置流管道,确认
channel()与handle()的绑定逻辑未被错误修改
- 检查处理器方法是否通过
Kafka容器重平衡或重启导致订阅关系丢失
当AWS MSK集群发生重平衡,或Spring Kafka容器重启时,可能出现通道订阅关系丢失:- 升级Spring Kafka和Spring Integration Kafka到兼容的稳定版本(当前使用的2.8.4/5.5.10可考虑升级到2.9.x/5.6.x系列)
- 配置容器的
ConsumerRebalanceListener,在重平衡完成后重新校验订阅者状态 - 确保通道配置
autoStartup=true,保证容器重启后通道自动恢复订阅
消息处理异常导致订阅者退订
如果处理器在处理消息时抛出未捕获的异常,可能导致订阅者被移除:- 优化处理器代码,确保所有异常都被捕获并处理,避免抛出未检查异常
- 配置Spring Kafka的错误处理机制,比如
ErrorHandler或DeadLetterPublishingRecoverer,将异常消息转发到死信队列,避免异常扩散影响订阅关系
内容的提问来源于stack exchange,提问作者marios390
相关产品推荐
相关产品推荐

