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:调整适配器配置参数
- 保留
cleanSession=false配置(已设置):该参数让Broker保留客户端的订阅信息,重连时自动恢复订阅; - 给适配器添加重连间隔配置:
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
相关产品推荐
相关产品推荐

