Spring Cloud Stream偶发MessageDispatchingException无订阅者问题求助
问题排查:Spring Cloud Stream Kafka消费停止,报Dispatcher无订阅者
场景概述
开发了一个Spring Cloud Stream应用,用于消费Kafka Topic并更新ElasticSearch索引,运行数天后出现异常,无法继续消费消息。
核心代码
@Bean public Consumer<Flux<Message<GraphTextKafkaRecord>>> fetchSeed() { return messages -> messages .map(message -> { var ack = (Acknowledgment) message.getHeaders().get(ACKNOWLEDGEMENT_KEY); var graphTextRecord = message.getPayload(); log.debug("Fetch message from kafka. message: {}", graphTextRecord); return Tuples.of(ack, graphTextRecord); }) .filter(tuple -> { var result = appRules.test(tuple.getT2()); if (!result) { tuple.getT1().acknowledge(); log.debug("Message: {} has been filtered due to rules defined in application", tuple.getT2()); } return result; }) .map(tuple -> { GraphTextKafkaRecord kafkaRecord = tuple.getT2(); return PageQuality .builder(tuple.getT1(), kafkaRecord.url(), kafkaRecord.pageQuality()) .pageQualityAlphaScore(kafkaRecord.pageQualityAlpha()) .statusCode(kafkaRecord.statusCode()) .build(); }) .flatMap(esHandler::insertPageQualityScore) .retryWhen(Retry.fixedDelay(5, Duration.ofSeconds(5))) .subscribe(); }
配置信息
spring: cloud: stream: default-binder: kafka kafka: binder: auto-create-topics: false brokers: ${PAGE_QUALITY_PROD_KAFKA_BROKERS:x.x.x.x:9092} enable-observation: false consumer-properties: allow.auto.create.topics: false bindings: fetchSeed-in-0: consumer: ack-mode: manual enable-dlq: false poll-timeout: 21474836470 start-offset: earliest bindings: fetchSeed-in-0: group: page-quality-group-prod destination: ${PAGE_QUALITY_PROD_KAFKA_TOPIC_NAME:graph-text} consumer: max-attempts: 10 back-off-initial-interval: 500 back-off-max-interval: 200 back-off-multiplier: 2.0 elasticsearch: uris: ${PAGE_QUALITY_PROD_ELASTIC_SEARCH_HOST:http://x.x.x.x:9200} username: ${PAGE_QUALITY_PROD_ELASTIC_SEARCH_USERNAME:user} password: ${PAGE_QUALITY_PROD_ELASTIC_SEARCH_PASSWORD:pass} socket-timeout: 30s connection-timeout: 30S
抛出异常
org.springframework.kafka.listener.ListenerExecutionFailedException: Listener failed at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.decorateException(KafkaMessageListenerContainer.java:2944) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2891) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeOnMessage(KafkaMessageListenerContainer.java:2857) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.lambda$doInvokeRecordListener$56(KafkaMessageListenerContainer.java:2780) at io.micrometer.observation.Observation.observe(Observation.java:559) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeRecordListener(KafkaMessageListenerContainer.java:2778) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeWithRecords(KafkaMessageListenerContainer.java:2630) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeRecordListener(KafkaMessageListenerContainer.java:2516) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:2168) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeIfHaveRecords(KafkaMessageListenerContainer.java:1523) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1487) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1362) at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(CompletableFuture.java:1804) at java.base/java.lang.Thread.run(Thread.java:833) Caused by: org.springframework.kafka.KafkaException: Failed to execute runnable at org.springframework.integration.kafka.inbound.KafkaInboundEndpoint.doWithRetry(KafkaInboundEndpoint.java:75) at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationRecordMessageListener.onMessage(KafkaMessageDrivenChannelAdapter.java:461) at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationRecordMessageListener.onMessage(KafkaMessageDrivenChannelAdapter.java:425) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2877) ... 12 more Caused by: org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'application.fetchSeed-in-0'., failedMessage=GenericMessage [payload=byte[88549], headers={kafka_offset=9184809, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@19b7a34e, deliveryAttempt=10, kafka_timestampType=CREATE_TIME, kafka_receivedPartitionId=34, kafka_receivedMessageKey=[B@2d72e8c7, kafka_receivedTopic=graph-text, kafka_receivedTimestamp=1678842937842, kafka_acknowledgment=Acknowledgment for graph-text-34@9184809, contentType=application/json, kafka_groupId=page-quality-group-prod}] at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:76) at org.springframework.integration.channel.AbstractMessageChannel.sendInternal(AbstractMessageChannel.java:373) at org.springframework.integration.channel.AbstractMessageChannel.sendWithMetrics(AbstractMessageChannel.java:344) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:324) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:297) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:187) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:166) at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:47) at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:109) at org.springframework.integration.endpoint.MessageProducerSupport.lambda$sendMessage$1(MessageProducerSupport.java:262) at io.micrometer.observation.Observation.observe(Observation.java:492) at org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:262) at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.sendMessageIfAny(KafkaMessageDrivenChannelAdapter.java:394) at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationRecordMessageListener.lambda$onMessage$0(KafkaMessageDrivenChannelAdapter.java:464) at org.springframework.integration.kafka.inbound.KafkaInboundEndpoint.lambda$doWithRetry$0(KafkaInboundEndpoint.java:70) at org.springframework.retry.support.RetryTemplate.doExecute(RetryTemplate.java:329) at org.springframework.retry.support.RetryTemplate.execute(RetryTemplate.java:225) at org.springframework.integration.kafka.inbound.KafkaInboundEndpoint.doWithRetry(KafkaInboundEndpoint.java:66) ... 15 more Caused by: org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[88549], headers={kafka_offset=9184809, kafka_consumer=org.apache.kafka.clients.consumer.KafkaConsumer@19b7a34e, deliveryAttempt=10, kafka_timestampType=CREATE_TIME, kafka_receivedPartitionId=34, kafka_receivedMessageKey=[B@2d72e8c7, kafka_receivedTopic=graph-text, kafka_receivedTimestamp=1678842937842, kafka_acknowledgment=Acknowledgment for graph-text-34@9184809, contentType=application/json, kafka_groupId=page-quality-group-prod}]
Actuator状态
{ "status": "UP", "components": { "binders": { "status": "UP", "components": { "kafka": { "status": "UP", "details": { "topicsInUse": [ "graph-text" ], "listenerContainers": [ { "isPaused": false, "listenerId": "KafkaConsumerDestination{consumerDestinationName='graph-text', partitions=0, dlqName='null'}.container", "isRunning": true, "groupId": "page-quality-group-prod", "isStoppedAbnormally": false } ] } } } } } }
排查方向
- 修复Flux订阅管理:代码中手动调用
.subscribe(),若流因异常终止会导致通道丢失订阅者。建议移除手动subscribe(),让Spring Cloud Stream框架自动管理订阅生命周期。 - 完善异常处理:当前仅在
flatMap阶段做重试,上游未捕获的异常会直接终止整个流。添加.onErrorContinue()或.onErrorResume()处理异常,避免流中断。 - 修正退避配置冲突:
back-off-max-interval设置为200ms小于初始间隔500ms,会导致退避策略异常,进而影响订阅状态。调整max-interval大于等于initial-interval。 - 规范手动ACK逻辑:在
filter阶段提前ACK,后续流处理异常会导致消费状态不一致。建议统一在处理完成后ACK,或依赖框架的自动ACK机制。 - 检查全链路日志:查看异常发生前后的应用日志,确认是否有流终止、Bean销毁或其他隐性异常,定位订阅者丢失的具体触发点。
内容的提问来源于stack exchange,提问作者Hatef Alipoor
相关产品推荐
相关产品推荐

