【问题标题】:spring cloud stream kafka: "Dispatcher has no subscribers" errorspring cloud stream kafka:“调度程序没有订阅者”错误
【发布时间】:2017-09-12 18:08:16
【问题描述】:

我正在使用 kafka binder 测试 spring cloud stream,但出现错误

原因:org.springframework.messaging.MessageDeliveryException:Dispatcher 没有频道“unknown.channel.name”的订阅者。;

pom.xml

<parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>1.4.0.RELEASE</version>
</parent>

<dependencies>
    <dependency>
        <groupId>org.springframework.cloud</groupId>
        <artifactId>spring-cloud-stream</artifactId>
        <version>1.1.2.RELEASE</version>
    </dependency>
    <dependency>
        <groupId>org.springframework.cloud</groupId>
        <artifactId>spring-cloud-stream-binder-kafka</artifactId>
        <version>1.1.2.RELEASE</version>
    </dependency>
</dependencies>

玩具应用程序是模仿秘书在员工和老板之间转移请求。

员工界面:

public interface SecretaryServingEmployee {
    @Output
    MessageChannel inbox();

    @Input
    SubscribableChannel rejected();

    @Input
    SubscribableChannel approved();

}

boss界面:

public interface SecretaryServingBoss {   
@Input
SubscribableChannel inbox();

@Output
MessageChannel rejected();

@Output
MessageChannel approved();      
}

application.properties

server.port=8080
spring.cloud.stream.bindings.inbox.destination=inbox
spring.cloud.stream.bindings.approved.destination=approved
spring.cloud.stream.bindings.rejected.destination=rejected

Employee.java

@EnableBinding(SecretaryServingEmployee.class)
@Component
public class Employee {
    private static Logger logger = LoggerFactory.getLogger(Employee.class);

    private SecretaryServingEmployee adminAssistent;

    @Autowired
    public Employee(SecretaryServingEmployee adminAssistent) {
        this.adminAssistent = adminAssistent;
    }

    @InboundChannelAdapter(value = "inbox")
    public String messageSource() {
        return "You are handsome!!";  // This is the message sent to boss
    }

    @ServiceActivator(inputChannel="approved")
    public void checkApproved(String message) {
        logger.info(":-)");
    }

    @ServiceActivator(inputChannel="rejected")
    public void checkRejected(RejectionLetter letter) {
        logger.warn(":-(");
    } 
}

Boss.java

@EnableBinding(SecretaryServingBoss.class)
@Component
public class Boss {
    private SecretaryServingBoss adminAssistent;

    @Autowired
    public Boss(SecretaryServingBoss adminAssistent) {
        this.adminAssistent = adminAssistent;
    }

    @ServiceActivator(inputChannel="inbox")
    public void sign(String content) {
        if (content.contains("You are handsome")) {
            adminAssistent.approved().send(message("nice work"));
        }
        else {
            adminAssistent.rejected().send(message("Don't send me shit"));
        }       
    }

    private <T> Message<T> message(T content) {
        return MessageBuilder.withPayload(content).build();
    }   

}

这是跟踪的一部分

org.springframework.messaging.MessageDeliveryException: Dispatcher has no subscribers for channel 'unknown.channel.name'.; nested exception is org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:81) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:423) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:373) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:115) ~[spring-messaging-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:45) ~[spring-messaging-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:105) ~[spring-messaging-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutput(AbstractMessageProducingHandler.java:292) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.handler.AbstractMessageProducingHandler.produceOutput(AbstractMessageProducingHandler.java:212) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.handler.AbstractMessageProducingHandler.sendOutputs(AbstractMessageProducingHandler.java:129) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:115) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:127) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.channel.FixedSubscriberChannel.send(FixedSubscriberChannel.java:70) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.channel.FixedSubscriberChannel.send(FixedSubscriberChannel.java:64) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:115) ~[spring-messaging-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.messaging.core.GenericMessagingTemplate.doSend(GenericMessagingTemplate.java:45) ~[spring-messaging-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.messaging.core.AbstractMessageSendingTemplate.send(AbstractMessageSendingTemplate.java:105) ~[spring-messaging-4.3.2.RELEASE.jar:4.3.2.RELEASE]
    at org.springframework.integration.endpoint.MessageProducerSupport.sendMessage(MessageProducerSupport.java:171) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.access$000(KafkaMessageDrivenChannelAdapter.java:47) ~[spring-integration-kafka-2.0.1.RELEASE.jar:na]
    at org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter$IntegrationMessageListener.onMessage(KafkaMessageDrivenChannelAdapter.java:197) ~[spring-integration-kafka-2.0.1.RELEASE.jar:na]
    at org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter$1.doWithRetry(RetryingAcknowledgingMessageListenerAdapter.java:76) ~[spring-kafka-1.0.5.RELEASE.jar:na]
    at org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter$1.doWithRetry(RetryingAcknowledgingMessageListenerAdapter.java:71) ~[spring-kafka-1.0.5.RELEASE.jar:na]
    at org.springframework.retry.support.RetryTemplate.doExecute(RetryTemplate.java:276) ~[spring-retry-1.1.3.RELEASE.jar:na]
    at org.springframework.retry.support.RetryTemplate.execute(RetryTemplate.java:172) ~[spring-retry-1.1.3.RELEASE.jar:na]
    at org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter.onMessage(RetryingAcknowledgingMessageListenerAdapter.java:71) ~[spring-kafka-1.0.5.RELEASE.jar:na]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.invokeListener(KafkaMessageListenerContainer.java:597) [spring-kafka-1.0.5.RELEASE.jar:na]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.access$1800(KafkaMessageListenerContainer.java:222) [spring-kafka-1.0.5.RELEASE.jar:na]
    at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer$ListenerInvoker.run(KafkaMessageListenerContainer.java:772) [spring-kafka-1.0.5.RELEASE.jar:na]
    at java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source) [na:1.8.0_121]
    at java.util.concurrent.FutureTask.run(Unknown Source) [na:1.8.0_121]
    at java.lang.Thread.run(Unknown Source) [na:1.8.0_121]
Caused by: org.springframework.integration.MessageDispatchingException: Dispatcher has no subscribers
    at org.springframework.integration.dispatcher.UnicastingDispatcher.doDispatch(UnicastingDispatcher.java:154) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.dispatcher.UnicastingDispatcher.dispatch(UnicastingDispatcher.java:121) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    at org.springframework.integration.channel.AbstractSubscribableChannel.doSend(AbstractSubscribableChannel.java:77) ~[spring-integration-core-4.3.1.RELEASE.jar:4.3.1.RELEASE]
    ... 29 common frames omitted

【问题讨论】:

    标签: spring-cloud-stream


    【解决方案1】:

    您的应用程序类似乎没有正确扫描组件。

    如果您将其作为 Spring Boot 应用程序运行,您能否确保要扫描的类是否已正确打包。

    例如,默认情况下,@SpringBootApplication 的组件扫描会查看 @SpringBootApplication 注释类所在的同一包下的类。

    【讨论】:

    • 是的,它是一个Spring Boot应用程序,所有代码都在一个包下。
    • 您不需要为每个目的地命名吗?例如@Input("bossInbox")
    • 唯一的目的地名称解决了这个问题。谢谢你,@GaryRussell
    【解决方案2】:

    我不时为我的 Kafka 制作人收到这个问题。在我的情况下,问题是 Kafka 生产者没有被 Spring 完全加载,我的代码试图使用该生产者对象发送一个 msg。一个简单的解决方法是在启动期间将消息发送延迟几秒钟,或者使用更复杂的解决方案,例如倒计时闩锁或弹簧事件。

    你想在你的代码开始向 kafka 生产者发送 msg 之前看到这样的一行:

    00:38:54.107 [Thread-4] INFO o.a.k.clients.producer.KafkaProducer - [Producer clientId=trade-data-producer-exasol-2] 实例化了一个幂等生产者。

    【讨论】:

      【解决方案3】:

      有人可能会从我的经验中受益,整天都在与这个问题作斗争,直到我看到@Ilayaperumal 的回复。

      我继承了一个 groovy 微服务,声明了组件扫描,想知道他们为什么这样做;

      例如,默认情况下,@SpringBootApplication 的组件扫描会查看 >@SpringBootApplication >注释类所在的同一包下的类。

      所以我在没有手动更新组件扫描值的情况下重构了包,这就是我的问题的根源。

      如果您遇到同样的问题,请检查您的包裹扫描。

      【讨论】:

        猜你喜欢
        • 2017-04-10
        • 2018-09-17
        • 2017-08-23
        • 2015-11-02
        • 2013-08-16
        • 2018-02-19
        • 2020-04-22
        • 1970-01-01
        • 2016-12-31
        相关资源
        最近更新 更多