【第二章16】深入理解MQTT订阅选项:从基础概念到实战应用
1. 引言为什么需要订阅选项MQTTMessage Queuing Telemetry Transport是一种轻量级的发布/订阅消息传输协议专为低带宽、高延迟或不稳定的网络环境设计。在物联网IoT、移动应用和微服务通信中MQTT 因其高效和可靠而广受欢迎。在 MQTT 通信模型中客户端通过**订阅Subscribe**来接收其感兴趣的主题Topic上的消息。然而简单的订阅有时无法满足复杂的业务需求例如如何只接收特定质量的消息如何避免收到过时的历史消息如何控制服务器为离线客户端保留消息的数量MQTT 订阅选项Subscription Options正是为了解决这些问题而设计的。它们允许客户端在订阅时指定一系列参数从而精细地控制消息的接收行为。本文将系统性地介绍 MQTT 订阅选项的基本概念、工作原理、实战代码以及典型应用场景。2. MQTT 订阅选项详解2.1 QoS服务质量等级QoS 是 MQTT 协议中保证消息可靠性的核心机制。在订阅时指定的 QoS 等级决定了客户端从服务器接收消息的最大努力保证级别。QoS 等级名称描述消息是否可能重复消息是否可能丢失0At most once (最多一次)发完即忘不保证送达。否是1At least once (至少一次)确保消息至少送达一次可能重复。是否2Exactly once (恰好一次)确保消息恰好送达一次开销最大。否否订阅 QoS 的最终生效规则客户端在订阅时请求一个 QoS 等级例如 QoS 1而发布者在发布消息时也会指定一个 QoS 等级例如 QoS 2。最终客户端实际接收消息的 QoS 等级是两者中的较小值。示例客户端以 QoS 1 订阅主题sensor/temperature发布者以 QoS 2 向该主题发布消息。则客户端最终将以 QoS 1 的保证级别收到该消息。这是 MQTT 协议的设计旨在避免服务器向能力不足的客户端发送其无法处理的高 QoS 消息。2.2 No LocalNo Local选项用于控制客户端是否接收自己发布的消息。当设置为true时客户端将不会收到由它自己发布到所订阅主题的消息。应用场景聊天室用户发送一条消息后不希望在自己的客户端界面里再看到一次来自服务器的相同消息回声。设备控制环一个设备发布状态后又订阅了该状态主题。设置No Local可以避免它处理自己发出的状态更新从而防止逻辑循环。2.3 Retain As PublishedRetain As Published选项控制服务器在转发消息时是否保持消息原有的保留Retain标志。保留消息Retained Message当发布者发布一条消息时可以设置retaintrue。服务器会将该消息保存在主题下后续任何新订阅该主题的客户端都会立即收到这条最新的保留消息。默认行为Retain As Published false服务器在向订阅者转发消息时无论原消息的retain标志是什么都会将其置为false。因此订阅者无法区分收到的是一条实时消息还是一条保留消息。设置 Retain As Published true服务器将保持原消息的retain标志不变。这使得订阅者可以知道消息的来源性质便于进行不同的业务逻辑处理。2.4 Retain HandlingRetain Handling选项是 MQTT 5.0 引入的新特性用于控制订阅建立时客户端是否希望接收该主题上已有的保留消息。它有三个可选值值含义0发送保留消息默认。订阅建立时立即发送该主题下的保留消息。1仅当订阅是新建时发送保留消息。如果客户端重复订阅同一个主题例如为了修改 QoS则不发送。2不发送保留消息。订阅建立时忽略所有保留消息只接收后续的实时消息。应用场景设备初始化一个新设备上线订阅其配置主题希望立即获取最新的配置使用值 0。避免重复处理客户端在断线重连后重新订阅不希望再次处理已经处理过的保留配置使用值 1 或 2。实时数据流只关心未来的温度数据不关心历史保留的最后一个温度值使用值 2。2.5 订阅选项对比汇总下表横向对比了四个核心订阅选项的关键特性方便读者快速查阅和选择选项MQTT 版本默认值主要作用典型应用场景注意事项QoS3.1.1 5.00最多一次控制消息传递的可靠性保证级别1. 配置下发QoS 1/22. 实时传感器数据QoS 03. 控制命令QoS 1/21. 最终生效 QoS min(订阅 QoS, 发布 QoS)2. QoS 2 开销最大确保恰好一次No Local5.0false控制客户端是否接收自己发布的消息1. 聊天室避免回声2. 设备控制环防止循环3. 自发布自订阅场景仅 MQTT 5.0 支持3.1.1 无此选项Retain As Published5.0false控制服务器转发时是否保持原消息的保留标志1. 需要区分实时消息与保留消息的业务2. 审计或日志场景需知消息来源仅 MQTT 5.0 支持设为 true 时订阅者能感知 retain 标志Retain Handling5.00发送保留消息控制订阅建立时是否接收已有的保留消息1. 设备初始化值 02. 避免重复处理值 13. 纯实时数据流值 2仅 MQTT 5.0 支持合理设置可避免不必要的保留消息处理使用建议MQTT 3.1.1 用户只能使用 QoS其他选项不可用。MQTT 5.0 用户可组合使用所有选项建议根据业务场景仔细配置。兼容性使用新选项时需确保 Broker 和客户端库均支持 MQTT 5.0。3. 实战代码示例下面我们使用 Java 的 Eclipse Paho 库来演示如何设置和使用这些订阅选项。示例将包含一个订阅者和一个发布者。3.1 环境准备首先确保你有一个 MQTT 服务器Broker在运行。可以使用公共的test.mosquitto.org或本地安装的 Mosquitto。使用 Maven 添加 Eclipse Paho 客户端依赖dependencygroupIdorg.eclipse.paho/groupIdartifactIdorg.eclipse.paho.client.mqttv5/artifactIdversion1.2.5/version/dependency或者使用 Gradleimplementationorg.eclipse.paho:org.eclipse.paho.client.mqttv5:1.2.53.2 订阅者代码包含订阅选项importorg.eclipse.paho.mqttv5.client.IMqttToken;importorg.eclipse.paho.mqttv5.client.MqttAsyncClient;importorg.eclipse.paho.mqttv5.client.MqttConnectionOptions;importorg.eclipse.paho.mqttv5.client.persist.MemoryPersistence;importorg.eclipse.paho.mqttv5.common.MqttException;importorg.eclipse.paho.mqttv5.common.MqttMessage;importorg.eclipse.paho.mqttv5.common.packet.MqttProperties;importorg.eclipse.paho.mqttv5.common.packet.UserProperty;importjava.util.concurrent.CountDownLatch;importjava.util.concurrent.TimeUnit;publicclassMqttSubscriberWithOptions{publicstaticvoidmain(String[]args){// 1. 配置连接参数Stringbrokertcp://test.mosquitto.org:1883;// MQTT Broker 地址StringclientIdJavaSubscriberDemo;// 客户端ID需唯一MemoryPersistencepersistencenewMemoryPersistence();// 持久化方式内存CountDownLatchlatchnewCountDownLatch(1);// 用于保持程序运行try{// 2. 创建 MQTT v5 异步客户端MqttAsyncClientclientnewMqttAsyncClient(broker,clientId,persistence);// 3. 设置连接选项MqttConnectionOptionsconnOptsnewMqttConnectionOptions();connOpts.setCleanStart(true);// 清除会话不保留之前的订阅状态connOpts.setAutomaticReconnect(true);// 启用自动重连// 4. 设置消息到达回调核心处理接收到的消息client.setCallback(neworg.eclipse.paho.mqttv5.client.MqttCallback(){Overridepublicvoiddisconnected(org.eclipse.paho.mqttv5.common.MqttExceptiondisconnectResponse){System.out.println(Disconnected: disconnectResponse.getMessage());}OverridepublicvoidmqttErrorOccurred(org.eclipse.paho.mqttv5.common.MqttExceptionexception){System.out.println(MQTT Error: exception.getMessage());}OverridepublicvoidmessageArrived(Stringtopic,MqttMessagemessage){// 当消息到达时触发此方法System.out.println(Topic: topic, QoS: message.getQos(), Retain: message.isRetained(), Payload: newString(message.getPayload()));}OverridepublicvoiddeliveryComplete(IMqttTokentoken){// 发布完成回调订阅者不需要实现}OverridepublicvoidconnectComplete(booleanreconnect,StringserverURI){System.out.println(Connected to broker: serverURI);}OverridepublicvoidauthPacketArrived(intreasonCode,MqttPropertiesproperties){// 认证包到达回调}});// 5. 连接 BrokerIMqttTokenconnectTokenclient.connect(connOpts);connectToken.waitForCompletion();// 等待连接完成System.out.println(Connected with result code: connectToken.getResponse().getReasonCode());// 6. 创建订阅选项MQTT v5 新特性MqttPropertiessubscriptionPropertiesnewMqttProperties();// 设置 No Local true不接收自己发布的消息subscriptionProperties.setNoLocal(true);// 设置 Retain As Published true保持原消息的保留标志subscriptionProperties.setRetainAsPublished(true);// 设置 Retain Handling 0订阅时发送保留消息subscriptionProperties.setRetainHandling(0);// 设置订阅标识符可选用于关联订阅subscriptionProperties.setSubscriptionIdentifier(1);// 7. 订阅主题并应用订阅选项// 参数说明主题过滤器 sensor/#QoS1订阅属性对象client.subscribe(sensor/#,1,subscriptionProperties);System.out.println(Subscribed to topic sensor/# with options.);// 8. 保持程序运行等待消息到达System.out.println(Waiting for messages... (Press CtrlC to exit));latch.await();// 阻塞主线程直到 latch.countDown() 被调用}catch(MqttException|InterruptedExceptione){e.printStackTrace();}}}###3.3发布者代码 javaimportorg.eclipse.paho.mqttv5.client.IMqttToken;importorg.eclipse.paho.mqttv5.client.MqttAsyncClient;importorg.eclipse.paho.mqttv5.client.MqttConnectionOptions;importorg.eclipse.paho.mqttv5.client.persist.MemoryPersistence;importorg.eclipse.paho.mqttv5.common.MqttException;importorg.eclipse.paho.mqttv5.common.MqttMessage;importcom.fasterxml.jackson.databind.ObjectMapper;importjava.util.concurrent.TimeUnit;publicclassMqttPublisherDemo{publicstaticvoidmain(String[]args){// 1. 配置连接参数Stringbrokertcp://test.mosquitto.org:1883;// MQTT Broker 地址StringclientIdJavaPublisherDemo;// 客户端ID需唯一MemoryPersistencepersistencenewMemoryPersistence();// 持久化方式ObjectMapperobjectMappernewObjectMapper();// JSON 序列化工具try{// 2. 创建 MQTT v5 异步客户端MqttAsyncClientclientnewMqttAsyncClient(broker,clientId,persistence);// 3. 设置连接选项MqttConnectionOptionsconnOptsnewMqttConnectionOptions();connOpts.setCleanStart(true);// 清除会话connOpts.setAutomaticReconnect(true);// 启用自动重连// 4. 连接 BrokerIMqttTokenconnectTokenclient.connect(connOpts);connectToken.waitForCompletion();// 等待连接完成System.out.println(Publisher connected to broker.);// 5. 准备传感器数据对象SensorDatasensorDatanewSensorData(25.5,60);// 6. 循环发布5条消息for(inti0;i5;i){// 6.1 发布普通温度数据QoS 1非保留消息StringtemperatureTopicsensor/temperature;StringtemperaturePayloadobjectMapper.writeValueAsString(sensorData);MqttMessagetemperatureMessagenewMqttMessage(temperaturePayload.getBytes());temperatureMessage.setQos(1);// 设置 QoS 等级为 1temperatureMessage.setRetained(false);// 非保留消息client.publish(temperatureTopic,temperatureMessage);System.out.println(Published normal message (i1): temperaturePayload);// 6.2 只在第一次发布保留配置消息QoS 1保留消息if(i0){StringconfigTopicsensor/config;StringconfigPayload{\interval\: 10};MqttMessageconfigMessagenewMqttMessage(configPayload.getBytes());configMessage.setQos(1);// 设置 QoS 等级为 1configMessage.setRetained(true);// 保留消息新订阅者会立即收到client.publish(configTopic,configMessage);System.out.println(Published retained config message: configPayload);}// 6.3 更新温度数据模拟传感器变化sensorData.setTemperature(sensorData.getTemperature()0.5);// 6.4 等待2秒模拟实际发布间隔TimeUnit.SECONDS.sleep(2);}// 7. 断开连接client.disconnect();System.out.println(Publisher disconnected.);}catch(MqttException|InterruptedException|com.fasterxml.jackson.core.JsonProcessingExceptione){e.printStackTrace();}}// 内部类传感器数据结构staticclassSensorData{privatedoubletemperature;privateinthumidity;publicSensorData(doubletemperature,inthumidity){this.temperaturetemperature;this.humidityhumidity;}publicdoublegetTemperature(){returntemperature;}publicvoidsetTemperature(doubletemperature){this.temperaturetemperature;}publicintgetHumidity(){returnhumidity;}publicvoidsetHumidity(inthumidity){this.humidityhumidity;}}}代码说明订阅者使用 MQTT v5 客户端在subscribe方法中通过MqttProperties对象设置了订阅选项No Local true不接收自己发布的消息Retain As Published true保持原消息的保留标志Retain Handling 0订阅时立即发送保留消息Subscription Identifier 1订阅标识符发布者发布两种消息普通的温度数据和一条保留的配置消息使用 Jackson 库进行 JSON 序列化。运行订阅者后再运行发布者。观察订阅者控制台输出可以看到订阅建立时立即收到了sensor/config主题的保留消息retaintrue。随后收到实时的sensor/temperature消息retainfalse。由于设置了No Local true如果同一个客户端也发布消息到sensor/#它将不会收到自己发布的消息。注意Eclipse Paho Java 客户端中设置订阅选项需要通过MqttProperties对象具体属性设置方法请参考官方文档。3.4 预期运行结果运行上述代码后订阅者控制台将输出类似以下内容Connected to broker: tcp://test.mosquitto.org:1883 Connected with result code: 0 Subscribed to topic sensor/# with options. Waiting for messages... (Press CtrlC to exit) # 订阅建立时立即收到的保留消息 Topic: sensor/config, QoS: 1, Retain: true, Payload: {interval: 10} # 随后收到的实时温度消息每2秒一条 Topic: sensor/temperature, QoS: 1, Retain: false, Payload: {temperature:25.5,humidity:60} Topic: sensor/temperature, QoS: 1, Retain: false, Payload: {temperature:26.0,humidity:60} Topic: sensor/temperature, QoS: 1, Retain: false, Payload: {temperature:26.5,humidity:60} Topic: sensor/temperature, QoS: 1, Retain: false, Payload: {temperature:27.0,humidity:60} Topic: sensor/temperature, QoS: 1, Retain: false, Payload: {temperature:27.5,humidity:60}结果分析保留消息处理由于订阅者设置了Retain Handling 0在订阅建立时立即收到了sensor/config主题的保留消息retaintrue。实时消息接收随后每2秒收到一条sensor/temperature主题的实时消息retainfalse。QoS 生效所有消息的 QoS 都是 1符合发布者和订阅者 QoS 的最小值规则。No Local 效果如果同一个客户端同时作为发布者和订阅者由于设置了No Local true它将不会收到自己发布的消息回声。4.1 场景一物联网设备配置下发需求设备上线后立即获取最新配置之后只接收配置变更。实现配置主题如device/{id}/config始终以保留消息形式发布。设备订阅时设置Retain Handling 0确保上线即收。设置Retain As Published true让设备能区分“初始配置”和“配置更新”。4.2 场景二实时数据仪表盘需求仪表盘只显示实时数据不显示历史快照。实现数据主题如sensor//data以非保留消息发布。仪表盘订阅时设置Retain Handling 2彻底忽略任何保留消息。结合No Local true防止后台数据推送服务收到自己的消息回声。4.3 场景三可靠命令控制需求向设备发送控制命令必须确保送达QoS 2但设备能力有限只支持 QoS 1。实现命令发布到cmd/{deviceId}QoS2。设备以 QoS1 订阅该主题。根据订阅 QoS 最终生效规则设备将以 QoS 1 收到命令。这平衡了可靠性与设备资源。可在业务层增加命令ID和确认机制弥补 QoS 1 可能重复的不足。4.4 最佳实践总结明确需求选择 QoS对配置、命令使用 QoS 1 或 2对高频传感数据使用 QoS 0。善用保留消息用于存储主题的“最后已知状态”便于新订阅者快速初始化。使用No Local避免循环任何可能发布并订阅同一主题的客户端都应考虑启用此选项。升级到 MQTT 5.0以充分利用Retain Handling等更精细的控制选项。主题规划清晰的主题层级如company/region/device/type/data配合订阅选项能构建出强大且清晰的消息路由体系。5. 总结MQTT 订阅选项提供了超越简单主题匹配的消息流控制能力。通过合理配置 QoS、No Local、Retain As Published 和 Retain Handling开发者可以构建出更健壮、更高效、更符合业务逻辑的物联网与消息驱动应用。理解这些选项的细微差别并在设计之初就将其纳入架构考虑是成为一名高级 MQTT 开发者的关键一步。希望本文能帮助你更好地驾驭 MQTT 协议打造出更出色的解决方案。