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

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

报错原因与解决方案

核心原因

  1. BridgeHandler 误用:BridgeHandler 默认需要明确指定输出通道,或者消息携带replyChannel头信息。当前从MQTT收到的消息没有该头,且Bridge未配置输出通道,导致无法解析消息的目标目的地。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 10:40:30