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

Spring Boot与MQTT集成:连接断开后自动重订阅配置问题

Spring MQTT Integration重连后自动订阅失效解决方案

问题背景

使用Spring MQTT Integration构建MQTT Broker订阅客户端,代码运行正常,但连接断开后重新建立连接时,无法自动恢复对目标Topic的订阅,仅能恢复连接状态,订阅丢失。

用户配置代码如下:

import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.core.MessageProducer;
import org.springframework.integration.mqtt.core.ConsumerStopAction;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
import org.springframework.integration.mqtt.support.MqttHeaders;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;

@Configuration
public class MqttBeans {

    Logger logger = LoggerFactory.getLogger(MqttBeans.class);
    
    
    @Bean
    public MqttConnectOptions mqttConnectOptions() {
        
        MqttConnectOptions options = new MqttConnectOptions();
        options.setServerURIs(new String[] { "ssl://URL:8883" });
        options.setUserName("ubidot_bridge");
        String pass = "PASS";
        options.setPassword(pass.toCharArray());
        options.setCleanSession(false);
        options.setAutomaticReconnect(true);
        
        options.setConnectionTimeout(30);
        options.setKeepAliveInterval(90);
        options.setMqttVersion(MqttConnectOptions.MQTT_VERSION_3_1_1);
        
        return options;
    }
    

    @Bean
    public MqttPahoClientFactory mqttClientFactory(MqttConnectOptions options) {

        DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
        
        factory.setConnectionOptions( options );
        factory.setConsumerStopAction(ConsumerStopAction.UNSUBSCRIBE_NEVER);
        logger.info("Reconnected to the broker");
        
        return factory;
    }

    @Bean
    public MessageChannel mqttInputChannel() {
        return new DirectChannel();
    }
    
    @Bean
    public MqttPahoMessageDrivenChannelAdapter mqttPahoMessageDrivenChannelAdapterConfig(MqttConnectOptions options) {
        
        MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter("ubidot_bridge_in",
                mqttClientFactory(options), "#");

        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(2);
        adapter.setOutputChannel(mqttInputChannel());
        logger.info("Setting up inbound channel");
        return adapter;
    }
    
    @Bean
    public MessageProducer inbound(MqttPahoMessageDrivenChannelAdapter adapter) {
        return adapter;
    }
    

    
    @Bean
    @ServiceActivator(inputChannel = "mqttInputChannel")
    public MessageHandler handler() {

        logger.info("Setting up msg receiver handler");

        return new MessageHandler() {

            @Override
            public void handleMessage(Message<?> message) throws MessagingException {

                String topic = message.getHeaders().get(MqttHeaders.RECEIVED_TOPIC).toString();

                logger.info("Msg received .. Topic: " + topic);
                logger.info("Payload " + message.getPayload());

                System.out.println();
            }

        };
    }

    @Bean
    public MessageChannel mqttOutboundChannel() {
        return new DirectChannel();
    }

    @Bean
    @ServiceActivator(inputChannel = "mqttOutboundChannel")
    public MessageHandler mqttOutbound( MqttConnectOptions options ) {

        // clientId is generated using a random number
        MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler("ubidot_bridge_out", mqttClientFactory(options));
        messageHandler.setAsync(true);
        messageHandler.setDefaultTopic("#");
        messageHandler.setDefaultRetained(false);
        return messageHandler;
    }
}

核心问题分析

  • 虽然开启了AutomaticReconnect,但MqttPahoMessageDrivenChannelAdapter默认不会在重连后自动触发订阅逻辑;
  • ConsumerStopAction.UNSUBSCRIBE_NEVER仅阻止客户端停止时取消订阅,无法覆盖重连场景的订阅恢复。

具体解决方案

方案1:自定义MqttCallback监听重连事件

通过实现MqttCallbackExtended监听连接完成事件,在重连成功时主动触发订阅:

@Bean
public MqttPahoMessageDrivenChannelAdapter mqttPahoMessageDrivenChannelAdapterConfig(MqttConnectOptions options) {
    MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter("ubidot_bridge_in",
            mqttClientFactory(options), "#");

    adapter.setCompletionTimeout(5000);
    adapter.setConverter(new DefaultPahoMessageConverter());
    adapter.setQos(2);
    adapter.setOutputChannel(mqttInputChannel());
    
    // 添加自定义回调监听重连
    adapter.setCallback(new MqttCallbackExtended() {
        @Override
        public void connectComplete(boolean reconnect, String serverURI) {
            if (reconnect) {
                logger.info("重连成功,重新订阅Topic");
                // 重新订阅目标Topic,可按需指定多个Topic
                adapter.subscribe("#");
            }
        }

        @Override
        public void connectionLost(Throwable cause) {
            logger.error("MQTT连接断开", cause);
        }

        @Override
        public void messageArrived(String topic, MqttMessage message) throws Exception {
            // 消息处理由适配器自动转发,此处可留空或添加日志
        }

        @Override
        public void deliveryComplete(IMqttDeliveryToken token) {
            // 入站适配器无需处理消息投递完成事件
        }
    });
    
    logger.info("Inbound Channel 初始化完成");
    return adapter;
}

方案2:调整适配器配置参数

  1. 保留cleanSession=false配置(已设置):该参数让Broker保留客户端的订阅信息,重连时自动恢复订阅;
  2. 给适配器添加重连间隔配置:
adapter.setRecoveryInterval(3000); // 每3秒尝试恢复连接与订阅

方案3:监听Spring Integration连接事件

通过Spring事件监听机制,捕获连接建立/失败事件,在连接建立后触发订阅:

@Component
public class MqttConnectionListener {

    private final MqttPahoMessageDrivenChannelAdapter adapter;
    private final Logger logger = LoggerFactory.getLogger(MqttConnectionListener.class);

    // 构造注入适配器实例
    public MqttConnectionListener(MqttPahoMessageDrivenChannelAdapter adapter) {
        this.adapter = adapter;
    }

    @EventListener
    public void onConnectionEstablished(MqttConnectionEstablishedEvent event) {
        if ("ubidot_bridge_in".equals(event.getClientId())) {
            logger.info("MQTT连接建立,确认订阅状态");
            // 强制触发订阅,确保生效
            adapter.subscribe("#");
        }
    }

    @EventListener
    public void onConnectionFailed(MqttConnectionFailedEvent event) {
        if ("ubidot_bridge_in".equals(event.getClientId())) {
            logger.error("MQTT连接失败,等待重连", event.getCause());
        }
    }
}

关键注意事项

  • 确保MQTT Broker支持cleanSession=false的订阅持久化,部分Broker需额外配置;
  • 订阅Topic可根据业务需求调整,subscribe方法支持传入多个Topic参数;
  • 完善日志输出,便于排查重连与订阅过程中的问题。

内容的提问来源于stack exchange,提问作者tjxob

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 15:50:31