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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 08:48:11