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

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进行测试,结果如下:

  1. 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
  1. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 09:08:02