1. 引言:为什么需要订阅选项?
MQTT(Message Queuing Telemetry Transport)是一种轻量级的发布/订阅消息传输协议,专为低带宽、高延迟或不稳定的网络环境设计。在物联网(IoT)、移动应用和微服务通信中,MQTT 因其高效和可靠而广受欢迎。
在 MQTT 通信模型中,客户端通过**订阅(Subscribe)**来接收其感兴趣的主题(Topic)上的消息。然而,简单的订阅有时无法满足复杂的业务需求,例如:
- 如何只接收特定质量的消息?
- 如何避免收到过时的历史消息?
- 如何控制服务器为离线客户端保留消息的数量?
MQTT 订阅选项(Subscription Options)正是为了解决这些问题而设计的。它们允许客户端在订阅时指定一系列参数,从而精细地控制消息的接收行为。本文将系统性地介绍 MQTT 订阅选项的基本概念、工作原理、实战代码以及典型应用场景。
2. MQTT 订阅选项详解
2.1 QoS(服务质量等级)
QoS 是 MQTT 协议中保证消息可靠性的核心机制。在订阅时指定的 QoS 等级,决定了客户端从服务器接收消息的最大努力保证级别。
| QoS 等级 | 名称 | 描述 | 消息是否可能重复 | 消息是否可能丢失 |
|---|---|---|---|---|
| 0 | At most once (最多一次) | 发完即忘,不保证送达。 | 否 | 是 |
| 1 | At least once (至少一次) | 确保消息至少送达一次,可能重复。 | 是 | 否 |
| 2 | Exactly once (恰好一次) | 确保消息恰好送达一次,开销最大。 | 否 | 否 |
订阅 QoS 的最终生效规则:
客户端在订阅时请求一个 QoS 等级(例如 QoS 1),而发布者在发布消息时也会指定一个 QoS 等级(例如 QoS 2)。最终客户端实际接收消息的 QoS 等级,是两者中的较小值。
示例:客户端以 QoS 1 订阅主题
sensor/temperature,发布者以 QoS 2 向该主题发布消息。则客户端最终将以 QoS 1 的保证级别收到该消息。这是 MQTT 协议的设计,旨在避免服务器向能力不足的客户端发送其无法处理的高 QoS 消息。
2.2 No Local
No Local选项用于控制客户端是否接收自己发布的消息。当设置为true时,客户端将不会收到由它自己发布到所订阅主题的消息。
应用场景:
- 聊天室:用户发送一条消息后,不希望在自己的客户端界面里再看到一次来自服务器的相同消息回声。
- 设备控制环:一个设备发布状态后,又订阅了该状态主题。设置
No Local可以避免它处理自己发出的状态更新,从而防止逻辑循环。
2.3 Retain As Published
Retain As Published选项控制服务器在转发消息时,是否保持消息原有的保留(Retain)标志。
- 保留消息(Retained Message):当发布者发布一条消息时,可以设置
retain=true。服务器会将该消息保存在主题下,后续任何新订阅该主题的客户端都会立即收到这条最新的保留消息。 - 默认行为(Retain As Published = false):服务器在向订阅者转发消息时,无论原消息的
retain标志是什么,都会将其置为false。因此,订阅者无法区分收到的是一条实时消息还是一条保留消息。 - 设置 Retain As Published = true:服务器将保持原消息的
retain标志不变。这使得订阅者可以知道消息的来源性质,便于进行不同的业务逻辑处理。
2.4 Retain Handling
Retain Handling选项是 MQTT 5.0 引入的新特性,用于控制订阅建立时,客户端是否希望接收该主题上已有的保留消息。它有三个可选值:
| 值 | 含义 |
|---|---|
| 0 | 发送保留消息(默认)。订阅建立时,立即发送该主题下的保留消息。 |
| 1 | 仅当订阅是新建时发送保留消息。如果客户端重复订阅同一个主题(例如为了修改 QoS),则不发送。 |
| 2 | 不发送保留消息。订阅建立时,忽略所有保留消息,只接收后续的实时消息。 |
应用场景:
- 设备初始化:一个新设备上线,订阅其配置主题,希望立即获取最新的配置(使用值 0)。
- 避免重复处理:客户端在断线重连后重新订阅,不希望再次处理已经处理过的保留配置(使用值 1 或 2)。
- 实时数据流:只关心未来的温度数据,不关心历史保留的最后一个温度值(使用值 2)。
2.5 订阅选项对比汇总
下表横向对比了四个核心订阅选项的关键特性,方便读者快速查阅和选择:
| 选项 | MQTT 版本 | 默认值 | 主要作用 | 典型应用场景 | 注意事项 |
|---|---|---|---|---|---|
| QoS | 3.1.1 & 5.0 | 0(最多一次) | 控制消息传递的可靠性保证级别 | 1. 配置下发(QoS 1/2) 2. 实时传感器数据(QoS 0) 3. 控制命令(QoS 1/2) | 1. 最终生效 QoS = min(订阅 QoS, 发布 QoS) 2. QoS 2 开销最大,确保恰好一次 |
| No Local | 5.0 | false | 控制客户端是否接收自己发布的消息 | 1. 聊天室避免回声 2. 设备控制环防止循环 3. 自发布自订阅场景 | 仅 MQTT 5.0 支持,3.1.1 无此选项 |
| Retain As Published | 5.0 | false | 控制服务器转发时是否保持原消息的保留标志 | 1. 需要区分实时消息与保留消息的业务 2. 审计或日志场景需知消息来源 | 仅 MQTT 5.0 支持;设为 true 时订阅者能感知 retain 标志 |
| Retain Handling | 5.0 | 0(发送保留消息) | 控制订阅建立时是否接收已有的保留消息 | 1. 设备初始化(值 0) 2. 避免重复处理(值 1) 3. 纯实时数据流(值 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 客户端依赖:
<dependency><groupId>org.eclipse.paho</groupId><artifactId>org.eclipse.paho.client.mqttv5</artifactId><version>1.2.5</version></dependency>或者使用 Gradle:
implementation'org.eclipse.paho:org.eclipse.paho.client.mqttv5:1.2.5'3.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. 配置连接参数Stringbroker="tcp://test.mosquitto.org:1883";// MQTT Broker 地址StringclientId="JavaSubscriberDemo";// 客户端ID,需唯一MemoryPersistencepersistence=newMemoryPersistence();// 持久化方式(内存)CountDownLatchlatch=newCountDownLatch(1);// 用于保持程序运行try{// 2. 创建 MQTT v5 异步客户端MqttAsyncClientclient=newMqttAsyncClient(broker,clientId,persistence);// 3. 设置连接选项MqttConnectionOptionsconnOpts=newMqttConnectionOptions();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. 连接 BrokerIMqttTokenconnectToken=client.connect(connOpts);connectToken.waitForCompletion();// 等待连接完成System.out.println("Connected with result code: "+connectToken.getResponse().getReasonCode());// 6. 创建订阅选项(MQTT v5 新特性)MqttPropertiessubscriptionProperties=newMqttProperties();// 设置 No Local = true(不接收自己发布的消息)subscriptionProperties.setNoLocal(true);// 设置 Retain As Published = true(保持原消息的保留标志)subscriptionProperties.setRetainAsPublished(true);// 设置 Retain Handling = 0(订阅时发送保留消息)subscriptionProperties.setRetainHandling(0);// 设置订阅标识符(可选,用于关联订阅)subscriptionProperties.setSubscriptionIdentifier(1);// 7. 订阅主题并应用订阅选项// 参数说明:主题过滤器 "sensor/#",QoS=1,订阅属性对象client.subscribe("sensor/#",1,subscriptionProperties);System.out.println("Subscribed to topic 'sensor/#' with options.");// 8. 保持程序运行,等待消息到达System.out.println("Waiting for messages... (Press Ctrl+C 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. 配置连接参数Stringbroker="tcp://test.mosquitto.org:1883";// MQTT Broker 地址StringclientId="JavaPublisherDemo";// 客户端ID,需唯一MemoryPersistencepersistence=newMemoryPersistence();// 持久化方式ObjectMapperobjectMapper=newObjectMapper();// JSON 序列化工具try{// 2. 创建 MQTT v5 异步客户端MqttAsyncClientclient=newMqttAsyncClient(broker,clientId,persistence);// 3. 设置连接选项MqttConnectionOptionsconnOpts=newMqttConnectionOptions();connOpts.setCleanStart(true);// 清除会话connOpts.setAutomaticReconnect(true);// 启用自动重连// 4. 连接 BrokerIMqttTokenconnectToken=client.connect(connOpts);connectToken.waitForCompletion();// 等待连接完成System.out.println("Publisher connected to broker.");// 5. 准备传感器数据对象SensorDatasensorData=newSensorData(25.5,60);// 6. 循环发布5条消息for(inti=0;i<5;i++){// 6.1 发布普通温度数据(QoS 1,非保留消息)StringtemperatureTopic="sensor/temperature";StringtemperaturePayload=objectMapper.writeValueAsString(sensorData);MqttMessagetemperatureMessage=newMqttMessage(temperaturePayload.getBytes());temperatureMessage.setQos(1);// 设置 QoS 等级为 1temperatureMessage.setRetained(false);// 非保留消息client.publish(temperatureTopic,temperatureMessage);System.out.println("Published normal message "+(i+1)+": "+temperaturePayload);// 6.2 只在第一次发布保留配置消息(QoS 1,保留消息)if(i==0){StringconfigTopic="sensor/config";StringconfigPayload="{\"interval\": 10}";MqttMessageconfigMessage=newMqttMessage(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.temperature=temperature;this.humidity=humidity;}publicdoublegetTemperature(){returntemperature;}publicvoidsetTemperature(doubletemperature){this.temperature=temperature;}publicintgetHumidity(){returnhumidity;}publicvoidsetHumidity(inthumidity){this.humidity=humidity;}}}代码说明:
- 订阅者使用 MQTT v5 客户端,在
subscribe方法中通过MqttProperties对象设置了订阅选项:No Local = true:不接收自己发布的消息Retain As Published = true:保持原消息的保留标志Retain Handling = 0:订阅时立即发送保留消息Subscription Identifier = 1:订阅标识符
- 发布者发布两种消息:普通的温度数据和一条保留的配置消息,使用 Jackson 库进行 JSON 序列化。
- 运行订阅者后,再运行发布者。观察订阅者控制台输出,可以看到:
- 订阅建立时,立即收到了
sensor/config主题的保留消息(retain=true)。 - 随后收到实时的
sensor/temperature消息(retain=false)。 - 由于设置了
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 Ctrl+C 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主题的保留消息(retain=true)。 - 实时消息接收:随后每2秒收到一条
sensor/temperature主题的实时消息(retain=false)。 - 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},QoS=2。 - 设备以 QoS=1 订阅该主题。
- 根据订阅 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 协议,打造出更出色的解决方案。