【问题标题】:Exception: org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel异常:org.springframework.messaging.MessageDeliveryException:调度程序没有频道订阅者
【发布时间】:2020-01-15 02:22:03
【问题描述】:

我有一个沙箱用于探索 Spring Cloud Stream 中新增的功能,但我在一个 Spring Cloud Stream 应用程序中使用 Function 和 Supplier 时遇到了问题。

在代码中,我使用了docs 中描述的示例。

首先,我在application.yml 中添加了相应的spring.cloud.stream.bindingsspring.cloud.stream.function.definition 属性到项目Function<String, String>。一切正常,我将消息发布到 my-fun-in Kafka 主题,应用程序执行功能并将结果发送到 my-fun-out 主题。

然后我将Supplier<Flux<String>> 添加到具有相应spring.cloud.stream.bindings 的同一个项目中,并将spring.cloud.stream.function.definition 值更新为fun;sup。这里奇怪的事情开始发生。当我尝试启动应用程序时,我收到以下错误:

2020-01-15 01:45:16.608 ERROR 10128 --- [oundedElastic-1] o.s.integration.handler.LoggingHandler   : org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'application.sup-out-0'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[20], headers={contentType=application/json, id=89301e00-b285-56e0-cb4d-8133555c8905, timestamp=1579045516603}], failedMessage=GenericMessage [payload=byte[20], headers={contentType=application/json, id=89301e00-b285-56e0-cb4d-8133555c8905, timestamp=1579045516603}]
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:77)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:453)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:403)
    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.router.AbstractMessageRouter.doSend(AbstractMessageRouter.java:206)
    at org.springframework.integration.router.AbstractMessageRouter.handleMessageInternal(AbstractMessageRouter.java:188)
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:170)
    at org.springframework.integration.handler.AbstractMessageHandler.onNext(AbstractMessageHandler.java:219)
    at org.springframework.integration.handler.AbstractMessageHandler.onNext(AbstractMessageHandler.java:57)
    at org.springframework.integration.endpoint.ReactiveStreamsConsumer$DelegatingSubscriber.hookOnNext(ReactiveStreamsConsumer.java:165)
    at org.springframework.integration.endpoint.ReactiveStreamsConsumer$DelegatingSubscriber.hookOnNext(ReactiveStreamsConsumer.java:148)
    at reactor.core.publisher.BaseSubscriber.onNext(BaseSubscriber.java:160)
    at reactor.core.publisher.FluxDoFinally$DoFinallySubscriber.onNext(FluxDoFinally.java:123)
    at reactor.core.publisher.EmitterProcessor.drain(EmitterProcessor.java:426)
    at reactor.core.publisher.EmitterProcessor.onNext(EmitterProcessor.java:268)
    at reactor.core.publisher.FluxCreate$BufferAsyncSink.drain(FluxCreate.java:793)
    at reactor.core.publisher.FluxCreate$BufferAsyncSink.next(FluxCreate.java:718)
    at reactor.core.publisher.FluxCreate$SerializedSink.next(FluxCreate.java:153)
    at org.springframework.integration.channel.FluxMessageChannel.doSend(FluxMessageChannel.java:63)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:453)
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:403)
    at org.springframework.integration.channel.FluxMessageChannel.lambda$subscribeTo$2(FluxMessageChannel.java:83)
    at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onNext(FluxPeekFuseable.java:189)
    at reactor.core.publisher.FluxPublishOn$PublishOnSubscriber.runAsync(FluxPublishOn.java:398)
    at reactor.core.publisher.FluxPublishOn$PublishOnSubscriber.run(FluxPublishOn.java:484)
    at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:84)
    at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:37)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
    at java.base/java.lang.Thread.run(Thread.java:834)
Caused by: org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers, failedMessage=GenericMessage [payload=byte[20], headers={contentType=application/json, id=89301e00-b285-56e0-cb4d-8133555c8905, timestamp=1579045516603}]
    at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:139)
    at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:106)
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:73)
    ... 34 more

之后我尝试了几件事:

  1. spring.cloud.stream.function.definition 还原为fun(禁用sup bean 绑定到外部目标)。应用程序启动,功能工作,供应商没有工作。一切都符合预期。
  2. spring.cloud.stream.function.definition 更改为sup(禁用fun bean 绑定到外部目标)。应用程序启动,功能不起作用,供应商工作(每秒向my-sup-out 主题生成消息)。一切都和预期一样。
  3. 已将 spring.cloud.stream.function.definition 值更新为 fun;sup。应用程序没有启动,得到同样的 MessageDeliveryException。
  4. spring.cloud.stream.function.definition 值交换为sup;fun。应用程序启动,供应商工作,但功能不起作用(没有向my-fun-out主题发送消息)。

最后一个比错误更让我困惑)所以现在我需要有人帮助解决问题。

我在配置中错过了什么吗?为什么在spring.cloud.stream.function.definition 中更改以; 分隔的bean 顺序会导致不同的结果?

完整项目上传至GitHub并添加如下:

StreamApplication.java:

package com.kaine;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import reactor.core.publisher.Flux;

import java.util.function.Function;
import java.util.function.Supplier;

@SpringBootApplication
public class StreamApplication {

    public static void main(String[] args) {
        SpringApplication.run(StreamApplication.class);
    }

    @Bean
    public Function<String, String> fun() {
        return value -> value.toUpperCase();
    }

    @Bean
    public Supplier<Flux<String>> sup() {
        return () -> Flux.from(emitter -> {
            while (true) {
                try {
                    emitter.onNext("Hello from Supplier!");
                    Thread.sleep(1000);
                } catch (Exception e) {
                    // ignore
                }
            }
        });
    }
}

application.yml

spring:
  cloud:
    stream:
      function:
        definition: fun;sup
      bindings:
        fun-in-0:
          destination: my-fun-in
        fun-out-0:
          destination: my-fun-out
        sup-out-0:
          destination: my-sup-out

build.gradle.kts:

plugins {
    java
}

group = "com.kaine"
version = "1.0-SNAPSHOT"

repositories {
    mavenCentral()
}

dependencies {
        implementation(platform("org.springframework.cloud:spring-cloud-dependencies:Hoxton.SR1"))
        implementation("org.springframework.cloud:spring-cloud-starter-stream-kafka")

        implementation(platform("org.springframework.boot:spring-boot-dependencies:2.2.2.RELEASE"))
}

configure<JavaPluginConvention> {
    sourceCompatibility = JavaVersion.VERSION_11
}

【问题讨论】:

    标签: java spring-cloud spring-cloud-stream spring-cloud-function


    【解决方案1】:

    实际上,这是我们文档的一个问题,因为我相信我们为他的案例提供了一个反应式供应商的坏例子。问题是您的供应商处于无限阻塞循环中。它基本上永远不会回来。 所以请将其更改为:

    @Bean
    public Supplier<Flux<String>> sup() {
        return () -> Flux.fromStream(Stream.generate(new Supplier<String>() {
    
            @Override
            public String get() {
                try {
                    Thread.sleep(1000);
                    return "Hello from Supplier";
                } catch (Exception e) {
                    // ignore
                }
            }
    
        })).subscribeOn(Schedulers.elastic()).share();
    }
    

    【讨论】:

    • Oleg,我已经尝试过您的代码 sn-p,它确实解决了我的问题。谢谢你的建议)小提示:get方法在try/catch之后需要return,否则编译器会大喊“缺少返回语句”
    • @Oleg 你能解释一下为什么 qestion 中的代码会阻塞吗?谢谢
    猜你喜欢
    • 2022-08-19
    • 2015-11-02
    • 1970-01-01
    • 2018-02-19
    • 2020-04-22
    • 1970-01-01
    • 1970-01-01
    • 2019-10-12
    • 2017-05-05
    相关资源
    最近更新 更多