Spring Integration MQTT报DestinationResolutionException异常求助
Spring Integration MQTT 报错:no output-channel or replyChannel header available 分析与解决
使用版本
- org.springframework.integration:spring-integration-mqtt:5.5.2
- org.springframework.boot:spring-boot-starter:2.5.3
- org.eclipse.paho:org.eclipse.paho.client.mqttv3:1.2.5
配置代码
@Configuration public class MqttConfig { @Bean public MqttPahoClientFactory mqttClientFactory() { DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory(); MqttConnectOptions options = new MqttConnectOptions(); options.setServerURIs(new String[] { "tcp://localhost:1883" }); return factory; } @Bean public MqttPahoMessageDrivenChannelAdapter inboundAdapter(MqttPahoClientFactory clientFactory) { return new MqttPahoMessageDrivenChannelAdapter("MyApp", clientFactory, "ReplyTopic"); } @Bean IntegrationFlow inboundFlow(MqttPahoMessageDrivenChannelAdapter inboundAdapter) { return IntegrationFlows.from(inboundAdapter) .bridge() .channel("replyChannel") .get(); } @Bean public MessageChannel replyChannel() { return MessageChannels.publishSubscribe().get();; } @Bean public MqttPahoMessageHandler outboundAdapter(MqttPahoClientFactory clientFactory) { return new MqttPahoMessageHandler("MyApp", clientFactory); } @Bean public IntegrationFlow outboundFlow(MqttPahoMessageHandler outboundAdapter) { return IntegrationFlows.from("requestChannel") .handle(outboundAdapter).get() } @MessagingGateway public interface MyGateway { @Gateway(requestChannel = "requestChannel", replyChannel = "replyChannel") String send(String request, @Header(MqttHeaders.TOPIC) String requestTopic); } }
客户端代码
@RestController public class MyController { @Autowired private MyGateway myGateway; @GetMapping("/sendRequest") public String sendRequest() { var response = myGateway.send("Hello", "MyTopic"); return response; } }
操作步骤与报错
执行curl http://localhost:8080/sendRequest发送请求,通过HiveMQ手动执行docker exec -it hivemq mqtt pub -t ReplyTopic -m "World" --debug回复消息,Spring应用输出报错:
2022-10-25 18:04:33.171 ERROR 17069 --- [T Call: MyApp] .m.i.MqttPahoMessageDrivenChannelAdapter : Unhandled exception for GenericMessage [payload=World, headers={mqtt_receivedRetained=false, mqtt_id=0, mqtt_duplicate=false, id=9dbd5e14-66ed-5dc8-6cea-6d04ef19c6cc, mqtt_receivedTopic=ReplyTopic, mqtt_receivedQos=0, timestamp=1666713873170}] org.springframework.messaging.MessageHandlingException: error occurred in message handler [org.springframework.integration.handler.BridgeHandler@6f63903c]; nested exception is org.springframework.messaging.core.DestinationResolutionException: no output-channel or replyChannel header available
报错原因与解决方案
核心原因
- BridgeHandler 误用:
BridgeHandler默认需要明确指定输出通道,或者消息携带replyChannel头信息。当前从MQTT收到的消息没有该头,且Bridge未配置输出通道,导致无法解析消息的目标目的地。 - MQTT客户端ID冲突:入站适配器和出站适配器使用了相同的
clientId("MyApp"),这会导致MQTT连接冲突,引发潜在的消息处理异常。
修复步骤
1. 移除多余的BridgeHandler
修改入站IntegrationFlow,直接将MQTT消息发送到replyChannel,无需Bridge中转:
@Bean IntegrationFlow inboundFlow(MqttPahoMessageDrivenChannelAdapter inboundAdapter) { return IntegrationFlows.from(inboundAdapter) .channel("replyChannel") .get(); }
2. 区分MQTT客户端ID
为入站和出站适配器配置不同的clientId,避免连接冲突:
// 入站适配器 @Bean public MqttPahoMessageDrivenChannelAdapter inboundAdapter(MqttPahoClientFactory clientFactory) { return new MqttPahoMessageDrivenChannelAdapter("MyApp-Inbound", clientFactory, "ReplyTopic"); } // 出站适配器 @Bean public MqttPahoMessageHandler outboundAdapter(MqttPahoClientFactory clientFactory) { return new MqttPahoMessageHandler("MyApp-Outbound", clientFactory); }
3. (可选)请求-回复关联优化
如果需要更严谨的请求-回复匹配,可以在发送请求时添加correlationId,并在回复消息中携带该ID,通过CorrelationStrategy和ReplyMessageCorrelator关联,但当前场景下,完成前两步即可解决报错问题。
内容的提问来源于stack exchange,提问作者Razalgul
相关产品推荐
相关产品推荐

