Spring MQTT调用toReactivePublisher转Flux报Dispatcher无订阅者如何解决
问题原因与解决方案
错误根因
你遇到的Dispatcher has no subscribers报错核心原因是以下两点:
- 你使用的
DirectChannel是点对点同步消息通道,只要通道没有注册订阅者,发送消息时就会直接抛出该异常。你仅将Flux声明为Bean并不会被Spring自动订阅,相当于通道始终没有消费者,MQTT适配器收到消息往通道发的时候自然触发报错。 - 你的原有写法存在冗余配置:你已经手动给MQTT入站适配器设置了
outputChannel为mqttInputChannel,再用IntegrationFlows.from(adapter)会覆盖适配器的输出通道配置,反而导致通道绑定逻辑混乱。
而你替换成@ServiceActivator注解的handler可以正常运行,是因为该注解会自动将handler注册为mqttInputChannel的订阅者,通道有了消费者就不会触发报错。
正确实现方案
方案1:使用响应式通道直接生成Flux(更推荐)
直接使用Spring Integration提供的FluxMessageChannel替代DirectChannel,原生适配响应式流场景:
- 替换通道实现
@Bean public MessageChannel mqttInputChannel() { // 替换为响应式消息通道,原生支持转Flux return new FluxMessageChannel(); }
- 保留原有MQTT适配器配置不变
@Bean public MessageProducerSupport mqttInbound() { MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter("SpringClient", mqttClientFactory(), topic()); adapter.setCompletionTimeout(10000); adapter.setConverter(new DefaultPahoMessageConverter()); adapter.setQos(1); adapter.setOutputChannel(mqttInputChannel()); return adapter; }
- 从通道生成Flux,必须主动触发订阅
@Bean public Flux<Message<byte[]>> mqttMessageFlux(FluxMessageChannel mqttInputChannel) { return Flux.from(mqttInputChannel) // 这里做你需要的消息转换逻辑 .map(msg -> { String processedPayload = msg.getPayload() + ", received from MQTT"; return MessageBuilder.withPayload(processedPayload.getBytes(StandardCharsets.UTF_8)) .copyHeaders(msg.getHeaders()) .build(); }) // 按需添加share()操作符,支持多订阅者消费同一份消息流 .share(); } // Spring启动完成后主动订阅Flux,才会真正开始消费消息 @EventListener(ApplicationReadyEvent.class) public void initMqttFluxSubscription(Flux<Message<byte[]>> mqttMessageFlux) { mqttMessageFlux.subscribe( msg -> log.info("收到MQTT消息:{}", new String(msg.getPayload(), StandardCharsets.UTF_8)), err -> log.error("MQTT消息流异常", err) ); }
方案2:用toReactivePublisher方法实现
如果你要沿用toReactivePublisher的写法,需要调整配置避免通道冲突:
- 去掉MQTT适配器的手动输出通道配置
@Bean public MessageProducerSupport mqttInbound() { MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter("SpringClient", mqttClientFactory(), topic()); adapter.setCompletionTimeout(10000); adapter.setConverter(new DefaultPahoMessageConverter()); adapter.setQos(1); // 删掉adapter.setOutputChannel(mqttInputChannel())这一行,交给IntegrationFlow管理 return adapter; }
- 生成Publisher再转Flux,同样需要主动订阅
@Bean public Flux<Message<byte[]>> mqttInFlow() { Publisher<Message<byte[]>> publisher = IntegrationFlows.from(mqttInbound()) .transform(p -> p + ", received from MQTT") .transform(String.class, s -> s.getBytes(StandardCharsets.UTF_8)) .toReactivePublisher(); return Flux.from(publisher).share(); } // 同样需要主动订阅这个Flux,和方案1的订阅逻辑完全一致 @EventListener(ApplicationReadyEvent.class) public void initMqttFluxSubscription(Flux<Message<byte[]>> mqttInFlow) { mqttInFlow.subscribe( msg -> log.info("收到MQTT消息:{}", new String(msg.getPayload(), StandardCharsets.UTF_8)), err -> log.error("MQTT消息流异常", err) ); }
内容的提问来源于stack exchange,提问作者florian1
相关产品推荐
相关产品推荐

