Paho MQTT v1.2.5多订阅匹配消息时重复接收问题咨询
Paho MQTT v1.2.5多订阅匹配消息时重复接收问题咨询
测试环境:Paho版本1.2.5
我最近遇到一个MQTT订阅的奇怪问题:我往主题root/msg/1/data发送一条消息,同时客户端订阅了两个主题——root/msg/1/#和root/msg/+/#。按道理这两个订阅都能匹配上发送的主题,我原本以为每个订阅会触发一次消息监听器,也就是总共收到两条消息。但实际运行后,每个订阅都让监听器触发了两次,总共收到四条消息,这和我的预期完全不符。
下面是我用到的测试代码:
public class MqttPahoExample { public static void main(String[] args) throws MqttException { MqttConnectionOptionsBuilder builder = new MqttConnectionOptionsBuilder(); MqttConnectionOptions mqttConnectionOptions = builder.automaticReconnect(true) .username("user") .password("password".getBytes()) .cleanStart(false) .requestResponseInfo(true) .build(); mqttConnectionOptions.setSendReasonMessages(false); MqttClient mqttClient = new MqttClient("tcp://localhost:61616", "Client-01"); mqttClient.connect(mqttConnectionOptions); // 两个订阅,都会匹配到发送的消息 MqttSubscription mqttSubscription1 = new MqttSubscription("root/msg/1/#", 2); MqttSubscription mqttSubscription2 = new MqttSubscription("root/msg/+/#", 2); IMqttMessageListener iMqttMessageListener = (topic, message1) -> System.out.println("topic:" + topic + " message:" + new String(message1.getPayload(), StandardCharsets.UTF_8) + " id:" + message1.getId()); MqttSubscription[] mqttSubscriptions = {mqttSubscription1, mqttSubscription2}; IMqttMessageListener[] mqttMessageListeners = {iMqttMessageListener, iMqttMessageListener}; mqttClient.subscribe(mqttSubscriptions, mqttMessageListeners); // 发布消息 MqttMessage message = new MqttMessage("TestMessage".getBytes()); message.setQos(2); mqttClient.publish("root/msg/1/data", message); } }
预期输出
topic:root/msg/1/data message:TestMessage id:1 topic:root/msg/1/data message:TestMessage id:1
实际输出
topic:root/msg/1/data message:TestMessage id:2 topic:root/msg/1/data message:TestMessage id:2 topic:root/msg/1/data message:TestMessage id:1 topic:root/msg/1/data message:TestMessage id:1
我原本预期每个订阅对应一条消息,结果却收到了双倍的数量。
补充测试(编辑内容)
一开始我用的是Paho MQTT5,后来根据建议尝试使用订阅标识符,但问题依然存在——每个重叠订阅还是会收到额外的消息副本。后来发现这个问题很大程度上取决于Broker对MQTT规范中这句话的实现:
Server MAY deliver further copies of the message, one for each additional matching subscription and respecting the subscription’s QoS in each case.
于是我换了不同的MQTT Broker进行测试,结果如下:
- Artemis Broker输出
root/msg/1/# topic:root/msg/1/data message:TestMessage id:1 root/msg/+/# topic:root/msg/1/data message:TestMessage id:1 root/msg/1/# topic:root/msg/1/data message:TestMessage id:2 root/msg/+/# topic:root/msg/1/data message:TestMessage id:2
- Mosquitto Broker输出
root/msg/1/# topic:root/msg/1/data message:TestMessage id:1 root/msg/+/# topic:root/msg/1/data message:TestMessage id:2
下面是我使用订阅标识符的测试代码:
public class MqttPahoFinalAsync { public static void main(String[] args) throws MqttException { MqttConnectionOptionsBuilder builder = new MqttConnectionOptionsBuilder(); MqttConnectionOptions mqttConnectionOptions = builder.automaticReconnect(true) .username("user") .password("password".getBytes()) .cleanStart(true) .requestReponseInfo(true) .build(); mqttConnectionOptions.setUseSubscriptionIdentifiers(true); MqttAsyncClient mqttClient = new MqttAsyncClient("tcp://localhost:1883", "Client-01"); mqttClient.connect(mqttConnectionOptions).waitForCompletion(); // 两个订阅,都会匹配到发送的消息 // 订阅1 MqttProperties subProperties1 = new MqttProperties(); subProperties1.setSubscriptionIdentifiers(List.of(0)); // Paho要求必须这样初始化 subProperties1.setSubscriptionIdentifier(1); MqttSubscription mqttSubscription1 = new MqttSubscription("root/msg/1/#", 2); IMqttMessageListener iMqttMessageListener1 = (topic, message1) -> System.out.println("root/msg/1/# topic:" + topic + " message:" + new String(message1.getPayload(), StandardCharsets.UTF_8) + " id:" + message1.getId()); mqttClient.subscribe(mqttSubscription1, null, null, iMqttMessageListener1, subProperties1).waitForCompletion(); // 订阅2 MqttProperties subProperties2 = new MqttProperties(); subProperties2.setSubscriptionIdentifiers(List.of(0)); // Paho要求必须这样初始化 subProperties2.setSubscriptionIdentifier(2); MqttSubscription mqttSubscription2 = new MqttSubscription("root/msg/+/#", 2); IMqttMessageListener iMqttMessageListener2 = (topic, message1) -> System.out.println("root/msg/+/# topic:" + topic + " message:" + new String(message1.getPayload(), StandardCharsets.UTF_8) + " id:" + message1.getId()); mqttClient.subscribe(mqttSubscription2, null, null, iMqttMessageListener2, subProperties2).waitForCompletion(); // 发布消息 MqttMessage message = new MqttMessage("TestMessage".getBytes()); message.setQos(2); mqttClient.publish("root/msg/1/data", message); } }
备注:内容来源于stack exchange,提问作者Arun Singh
相关产品推荐
相关产品推荐

