Spring Integration MQTT v5共享订阅无法正常工作问题求助
MQTT v5共享订阅无法接收消息问题排查与解决
问题描述
通过Spring Integration实现MQTT v5共享订阅,订阅主题为$share/group1/test,但共享订阅无法正常接收消息,仅非共享订阅可正常工作。
代码实现
@Bean public ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager() { MqttConnectionOptions connectionOptions = new MqttConnectionOptions(); connectionOptions.setServerURIs(new String[]{mqttProperties.getUrl()}); connectionOptions.setConnectionTimeout(30); connectionOptions.setMaxReconnectDelay(30); connectionOptions.setUserName(mqttProperties.getUsername()); connectionOptions.setPassword(mqttProperties .getPassword() .getBytes(StandardCharsets.UTF_8) ); connectionOptions.setAutomaticReconnect(true); Mqttv5ClientManager clientManager = new Mqttv5ClientManager( connectionOptions, Objects.requireNonNull(MyUtils.getClientId())); clientManager.setPersistence(new MemoryPersistence()); return clientManager; } @Bean public IntegrationFlow mqttInFlow( final ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager) { Mqttv5PahoMessageDrivenChannelAdapter messageProducer = new Mqttv5PahoMessageDrivenChannelAdapter( clientManager, "$share/group1/test"); messageProducer.setPayloadType(byte[].class); messageProducer.setManualAcks(true); messageProducer.setMessageConverter(new ByteArrayMessageConverter()); return IntegrationFlow.from(messageProducer) .handle(mqttMessageHandler) .get(); } public class MqttMessageHandler implements GenericHandler<byte[]> { @Override public Object handle(final byte[] payload, final MessageHeaders headers) { log.info("Received message: {}", payload); } }
已知条件
代码中订阅的是共享主题$share/group1/test
预期效果
发布到test主题的消息应能被MqttMessageHandler的handle方法接收并处理
问题现状
共享订阅无法正常接收消息,但非共享订阅可正常工作
问题分析与解决
核心原因
直接将$share/group1/test作为主题字符串传入Mqttv5PahoMessageDrivenChannelAdapter时,Spring Integration会对主题进行额外处理,导致实际发送给MQTT Broker的订阅格式不符合MQTT v5共享订阅的规范要求。
解决方案
通过Mqttv5Subscription对象配置共享订阅,而非手动拼接主题字符串,修改mqttInFlow实现:
@Bean public IntegrationFlow mqttInFlow( final ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager) { // 创建订阅配置,指定原始主题 Mqttv5Subscription subscription = new Mqttv5Subscription("test"); // 设置共享订阅组名 subscription.setSharedSubscriptionName("group1"); // 根据实际需求设置QoS级别 subscription.setQos(1); Mqttv5PahoMessageDrivenChannelAdapter messageProducer = new Mqttv5PahoMessageDrivenChannelAdapter(clientManager); // 添加配置好的共享订阅 messageProducer.addSubscription(subscription); messageProducer.setPayloadType(byte[].class); messageProducer.setManualAcks(true); messageProducer.setMessageConverter(new ByteArrayMessageConverter()); return IntegrationFlow.from(messageProducer) .handle(mqttMessageHandler) .get(); }
补充说明
- 使用
Mqttv5Subscription并设置sharedSubscriptionName后,Spring Integration会自动按照MQTT v5规范构造正确的共享订阅主题格式($share/group1/test),无需手动拼接。 - 需确保使用的MQTT Broker支持MQTT v5协议及共享订阅功能(如EMQX、Mosquitto 2.0+等)。
- 若开启了手动确认(
setManualAcks(true)),需在MqttMessageHandler的handle方法中手动确认消息,避免消息堆积或重复:
@Override public Object handle(final byte[] payload, final MessageHeaders headers) { log.info("Received message: {}", new String(payload, StandardCharsets.UTF_8)); // 手动确认消息 MqttMessageAcknowledgment acknowledgment = headers.get(MqttHeaders.MESSAGE_ACKNOWLEDGMENT, MqttMessageAcknowledgment.class); if (acknowledgment != null) { acknowledgment.acknowledge(); } return null; }
内容的提问来源于stack exchange,提问作者Aditya
相关产品推荐
相关产品推荐

