MQTT协议

分析 IoT、车机或者边缘网关流量时,MQTT 经常以一种很“好懂”的样子出现:主题像路径,负载像 JSON,抓包里还能直接看到 PUBLISH。真正麻烦的部分在连接断掉以后。同一条业务消息可能换一个 Packet Identifier 再出现;Broker 可以恢复订阅,却未必能跨进程重启恢复;DUP=0 也不能证明应用从未见过这条数据。

这篇文章从字节流开始,把 MQTT 3.1.1 的分帧、CONNECT、QoS 1/2、Session、Retained Message、Will 串成一套可以实际验证的模型。后半部分再单独分析 MQTT 5.0 的 Properties、Clean Start、Receive Maximum、Subscription Options 和 Shared Subscription。

文中的在线实验只绑定 127.0.0.1:18891,Broker 为 Mosquitto 2.0.18,客户端脚本只使用 Python 标准库并直接构造 MQTT 报文。MQTT 5 的十六进制报文是明确标注的合成样例,用于校验字段长度,没有伪装成网络抓包。

先建立一张协议地图

MQTT 是 Client/Server 发布订阅协议。客户端主动连接 Broker;一个客户端可以同时发布和订阅;Broker 保存订阅关系,并把 Application Message 转发给匹配的会话。MQTT 只把 payload 视为字节串,JSON、Protobuf 或图片格式都属于应用约定。

传感器、Broker、应用之间的发布订阅关系

一次业务发布通常经过两段 MQTT 交换:

1
publisher  -- MQTT exchange A -->  broker  -- MQTT exchange B -->  subscriber

这两段可以使用不同 QoS、不同 Packet Identifier 和不同重传状态。发布者收到 Broker 的 PUBACK,只能确认第一段完成;订阅者是否在线、是否处理成功,要由第二段和应用协议回答。MQTT 3.1.1 §4.3给出的 QoS 流程也是按一对发送方/接收方描述的。

MQTT 最小控制报文确实只有两个字节,例如:

1
2
3
C0 00    PINGREQ
D0 00 PINGRESP
E0 00 DISCONNECT(MQTT 3.1.1)

这只是 MQTT 协议长度,没有算 TCP、IP、TLS 和链路层开销。带主题、属性和 payload 的 PUBLISH 自然会长得多。

TCP 字节流怎样切出 MQTT 报文

固定头的“一字节加变长整数”

每个 MQTT Control Packet 先出现一个固定头首字节。高四位是报文类型,低四位是该类型的 flags。紧随其后的是 Remaining Length,它表示“固定头之后还有多少字节”,编码为一至四字节的 base-128 变长整数。

每个长度字节的低七位贡献数值,最高位说明后面是否还有长度字节:

1
value = Σ ((byte[i] & 0x7f) × 128^i)
Remaining Length 编码
0 00
127 7f
128 80 01
16,383 ff 7f
16,384 80 80 01
268,435,455 ff ff ff 7f

四个字节已经是协议上限。第四个长度字节仍带 continuation bit 时,解析器面对的是畸形报文,而不是“继续读第五字节”。规范的字段定义和最大值见 MQTT 3.1.1 §2.2.3 Remaining Length。

TCP 会拆,也会合

TCP 提供有序字节流,不保留应用报文边界。下面两个情况都很正常:

  • 第一次 recv() 只得到固定头和 Remaining Length 的第一个字节;
  • 一次 recv() 同时得到 PINGREQ || DISCONNECT 两个完整 MQTT 报文。

TCP 拆分、合并与 MQTT 增量解析

实验中的合成 PUBLISH 有 128 字节正文,因此开头是 30 80 01。脚本把 131 字节完整报文按 2 / 15 / 114 拆成三次 feed();前两次不能产出报文,第三次才返回一个完整 packet。另一个输入 c0 00 e0 00 则必须一次解析出两个 packet。

解析循环至少需要三个限制:

  1. Remaining Length 只能占一至四字节;
  2. fixed_header + RL_bytes + body 全部到齐后才交给报文解析器;
  3. 在分配大缓冲区之前先检查本地 maximum_packet_size。

核心逻辑可以缩成这样:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
def feed(chunk):
buffer.extend(chunk)
packets = []
while len(buffer) >= 2:
value, multiplier, end = 0, 1, None
for index in range(1, min(5, len(buffer))):
digit = buffer[index]
value += (digit & 0x7f) * multiplier
if not digit & 0x80:
end = index + 1
break
multiplier *= 128
if end is None:
if len(buffer) >= 5:
raise ValueError("malformed Remaining Length")
break
total = end + value
if total > MAX_PACKET_SIZE:
raise ValueError("packet too large")
if len(buffer) < total:
break
packets.append(bytes(buffer[:total]))
del buffer[:total]
return packets

真实客户端还要限制累计缓冲区、连接空闲时间和单连接并发状态。只检查 Remaining Length 的理论上限仍然可能允许对端声明 256 MiB 级报文。

CONNECT:连接参数都在这里定下来

一条最小的 MQTT 3.1.1 CONNECT

客户端 lab-a、Clean Session=1、Keep Alive=30 秒的报文可以写成:

1
2
3
4
5
6
10 11
00 04 4d 51 54 54 # Protocol Name = "MQTT"
04 # Protocol Level = 4,即 MQTT 3.1.1
02 # Connect Flags:Clean Session=1
00 1e # Keep Alive = 30
00 05 6c 61 62 2d 61

10 是 CONNECT 固定头,11 表示后面有 17 字节。本文客户端先等成功 CONNACK,再继续实验。注意这是一种明确处理连接结果的实现选择:MQTT 3.1.1 允许客户端在发出 CONNECT 后、不等 CONNACK 就继续发送控制报文;若 CONNECT 被拒绝,服务器不得处理其后的数据。服务器发给客户端的首个报文则必须是 CONNACK。MQTT 3.1.1 §3.1.4 与 §3.2分别规定了这两条方向不同的规则。

Connect Flags 不是随意组合的位图

位 字段 约束
7 User Name Flag 为 1 时 payload 中包含用户名字段
6 Password Flag 为 1 时包含密码字段;3.1.1 下用户名位为 0 时它也必须为 0
5 Will Retain Will Flag=0 时必须为 0
4–3 Will QoS 只能为 0、1、2;Will Flag=0 时必须为 0
2 Will Flag 在 Broker 上登记 Will Message
1 Clean Session 决定是否丢弃旧会话并在断开后丢弃本次会话
0 Reserved 必须为 0

CONNECT payload 的字段顺序由这些 flags 决定:Client ID、Will Topic、Will Message、User Name、Password。抓包解析器如果不先验证 flags,就可能从错误偏移读取后续字段。

MQTT UTF-8 字符串还有协议约束

MQTT 的 UTF-8 Encoded String 以两字节大端长度开头,长度是编码后的字节数,不是字符数。MQTT 3.1.1 强制禁止 U+0000、surrogate 区间和不合法 UTF-8;对列出的控制字符及 Unicode 非字符,使用的是 SHOULD NOT,接收方可以关闭连接。不能把这些推荐限制全部写成相同的强制禁用规则。MQTT 3.1.1 §1.5.3列出了具体约束。

因此,一个稳妥的解码器需要依次检查:

1
2
3
4
5
长度前缀是否完整
→ 字节数是否在当前报文范围内
→ UTF-8 编码是否合法
→ MQTT 禁止的码点是否出现
→ 再交给主题、Client ID 等字段自己的规则

用户名和密码在 MQTT 3.1.1 中并不具有相同的数据类型:User Name 是 UTF-8 Encoded String,Password 是带两字节长度的二进制数据。把密码强制按 UTF-8 解码会改变协议允许的输入空间。

Client ID 标识会话,不自动证明身份

Client ID 用来把网络连接关联到 Session。两个客户端使用同一个 Client ID 时,Broker 会关闭原来的连接;现场出现周期性掉线时,Client ID 冲突应进入排查列表。

认证结果来自 Broker 的认证配置、证书或插件。若 ACL 使用 %c 绑定 Client ID,那是部署策略主动赋予它授权意义,仍需保证攻击者不能随意抢用这个 ID。

Packet Identifier 的所有权与复用

Packet Identifier 是 16 位非零整数。它出现在 QoS 1/2 PUBLISH 和 SUBSCRIBE、UNSUBSCRIBE 及相关确认中;QoS 0 PUBLISH 没有这个字段。

发送方不能在前一个交换完成前,为另一条需要 Packet Identifier 的未完成交换复用同一个标识符。相应交换完成后,它可以立即重用。规范在 MQTT 3.1.1 §2.3.1定义了“in use”和释放条件。

这个编号有三个边界:

  • 属于某个 Client—Server Session 中一个方向上的未完成协议交换;
  • 恢复同一持久 Session 时,未完成 PUBLISH/PUBREL 继续使用原编号;
  • 两个传输方向拥有各自的编号空间;
  • Broker 转发时会为下游交换分配自己的编号。

例如发布者到 Broker 使用 0x1234,Broker 发给订阅者时完全可以使用 0x0001。拿 Packet Identifier 做跨设备、跨重连的订单号,会在编号复用后制造错误去重。

后面的回环实验先让 0x1234 完成 QoS 2 的 PUBCOMP,再用同一个编号发送 payload after-reuse。Mosquitto 接受了第二个新交换,说明编号生命周期已经结束。

Topic Name、Topic Filter 与订阅结果

PUBLISH 携带 Topic Name;SUBSCRIBE 携带 Topic Filter。Topic Name 不能含 + 或 #,Filter 才能使用通配符。

Filter 匹配 不匹配
sensor/+/temperature sensor/10/temperature sensor/10/room/temperature
sensor/# sensor、sensor/10/temperature sensors/10
+/status car/status car/1/status

+ 占一个完整层级;# 覆盖剩余层级并且必须放在末尾。空层级同样参与匹配,所以 a//b 与 a/b 是两个不同主题。以 $ 开头的主题有单独规则:不以 $ 开头的通配符 Filter 不会匹配 $SYS/...。MQTT 3.1.1 §4.7给出了精确匹配要求。

SUBSCRIBE 中写的 QoS 是请求上限。Broker 在 SUBACK 中逐项返回授予的 QoS 或 0x80 失败码。后续转发等级不会超过发布消息的 QoS 与订阅授予值中较低的一方。

主题树最好同时表达租户和资源边界:

1
2
3
devices/42/telemetry/temperature
devices/42/status/online
devices/42/commands/reboot

应用不能只从 payload 的 device_id 决定权限。Broker 已认证为设备 42 的连接,应在 ACL 层就无法发布 devices/43/...。

QoS 0、1、2 到底保证什么

QoS 1 和 QoS 2 的确认时序

QoS 单段交换 发送端需要保存到何时 接收语义
0 PUBLISH 写入网络后无需协议状态 至多一次,可能丢失
1 PUBLISH → PUBACK 收到匹配 PUBACK 至少一次,可能重复
2 PUBLISH → PUBREC → PUBREL → PUBCOMP 收到匹配 PUBCOMP 该协议交换一次交付

QoS 0:没有协议确认

QoS 0 的 PUBLISH 不含 Packet Identifier,也没有 PUBACK。TCP 成功写入只表示数据进入本端网络栈;连接在传输途中失败时,应用无法从 MQTT 确认对端是否接收。

它适合下一条数据很快覆盖上一条的高频遥测,也可能适合应用本身已有批次确认的场景。是否可以丢,取决于业务数据而不是“设备性能弱”这个标签。

QoS 1:PUBACK 丢失会制造重复

发送方保存未确认的 PUBLISH。连接恢复并继续同一 Session 时,未确认的 QoS 1 PUBLISH 要以 DUP=1 重发;接收方可能已经处理过第一次发送,只是 PUBACK 没能返回。

1
2
3
4
5
6
7
8
9
10
11
12
Client                         Broker
| PUBLISH d0 q1 id=7 |
|----------------------------->| Broker 接收并可能转发
| PUBACK id=7(丢失) |
|<--------------X--------------|
| 网络断开 |
| CONNECT clean=0 |
|----------------------------->|
| PUBLISH d1 q1 id=7 |
|----------------------------->| 同一 MQTT packet 的重传
| PUBACK id=7 |
|<-----------------------------|

DUP=1 表示当前发送方可能在重发这个 PUBLISH Packet。Broker 把 Application Message 转发给另一个连接时,不复制上游 DUP 位,而是按下游交换的状态设置。因此订阅端可能看到 DUP=0 的业务重复:例如发布者重建了一个全新交换,或应用主动重发并获得了新编号。规范在 MQTT 3.1.1 §3.3.1.1明确指出 DUP 不会由 Broker 原样传播。

MQTT 3.1.1 的重传要求与 Session 重连关联。实现不应仅凭一个任意的应用定时器,在仍然连接时反复发送相同未确认 packet;库通常会根据其协议版本和状态机处理恢复。

QoS 2:四步交换消除一次协议事务内的重复交付

QoS 2 的接收方在第一次 PUBLISH 后取得消息所有权、保存足够的去重状态并回复 PUBREC。发送方收到 PUBREC 后可以丢弃原始 Application Message,但要保存 PUBREL 状态。规范允许接收方采用两种方法:保存完整消息直到 PUBREL,或较早向后续接收者推进并保存 Packet Identifier;两种方法都必须保证同一交换不会因重复 PUBLISH 再次向后投递。收到 PUBREL 后,接收方释放对应状态并回复 PUBCOMP。

若 PUBLISH/PUBREC 阶段断线,恢复同一 Session 后发送方重发 PUBLISH DUP=1。若 PUBREL/PUBCOMP 阶段断线,则重发 PUBREL。两条路径都依赖双方保留状态。

QoS 2 两段独立状态机与重复 PUBLISH

回环脚本故意执行下面的上游序列:

1
2
3
4
5
6
34 ... 12 34 ...    PUBLISH, DUP=0, QoS=2, PID=0x1234
50 02 12 34 PUBREC
3c ... 12 34 ... 同内容 PUBLISH, DUP=1, QoS=2, PID=0x1234
50 02 12 34 PUBREC
62 02 12 34 PUBREL
70 02 12 34 PUBCOMP

Mosquitto 对重复 PUBLISH 再次回复 PUBREC,但订阅端只收到一次 payload。Broker 转发时使用自己的出站 PID,并与订阅端完成另一条 QoS 2 状态机。MQTT 3.1.1 §4.3.3描述了两种允许的接收端处理方法及其状态要求。

QoS 2 仍然不是业务事务

假设设备收到 commands/open 后做三步:驱动继电器、写数据库、发布结果。MQTT QoS 2 约束的是某一段协议交换;以下情况仍会让业务动作重复:

  • 生产者以新 Packet Identifier 再发布一次同一个命令;
  • Broker 集群、桥接或应用网关重新生成一条消息;
  • 消费进程完成外部动作后,在记录处理结果之前崩溃;
  • 操作员或补偿任务主动重发历史请求。

需要“只扣款一次”或“只开锁一次”时,应用要有稳定的业务标识和原子去重记录,不能拿 QoS 2 的四个报文代替数据库事务。

DUP 与 RETAIN 的作用域

原先的逐字节例子把 23.5 以 QoS 1、RETAIN=1 发布到 lab/temp:

1
33 10 00 08 6c 61 62 2f 74 65 6d 70 00 01 32 33 2e 35

MQTT 3.1.1 PUBLISH 的固定头、主题、Packet Identifier 与负载

字节 含义
33 PUBLISH;DUP=0、QoS=1、RETAIN=1
10 Remaining Length=16
00 08 Topic Name 的 UTF-8 字节数
6c ... 70 lab/temp
00 01 Packet Identifier=1
32 33 2e 35 payload 23.5

客户端发给 Broker 的 RETAIN=1 要求 Broker 更新该主题的保留存储。向已有订阅转发这次实时发布时,MQTT 3.1.1 的 Broker 把 RETAIN 设为 0;以后新建的匹配订阅收到保留副本时,RETAIN 才是 1。

这正是回环实验的结果:

1
retain scope: established=0 new-subscription=1

向同一主题发送 RETAIN=1 且 payload 长度为 0 的 PUBLISH,会删除保留消息。零长度 payload 在这里是协议操作,不等同于存储一条普通空值历史记录。完整动作见 MQTT 3.1.1 §3.3.1.3。

Connection、Session、持久化、Retained 与 Will

五个概念各自保存什么

概念 生命周期与内容
Network Connection 一条当前 TCP/TLS/WebSocket 连接
Session 订阅、未完成 QoS 状态、符合规则的待交付消息
Broker persistence Broker 是否把内存状态落盘并在自身重启后恢复
Retained Message 某 Topic Name 的最后一条保留消息
Will CONNECT 时登记、在异常断开条件下由 Broker 发布的消息

MQTT 3.1.1 使用 Clean Session:

  • Clean Session=1:连接建立时丢弃与 Client ID 关联的旧 Session;本次连接结束后也不保留;
  • Clean Session=0:请求恢复或创建可跨网络连接存在的 Session;
  • CONNACK 的 Session Present 告诉客户端是否恢复了已有 Session。

当客户端预期恢复会话,却收到 Session Present=0,应重建订阅和本地未完成状态。盲目假设 Broker 一定保存过订阅,会造成“连接成功但收不到消息”。Session 状态的规范列表见 MQTT 3.1.1 §4.1。协议会话跨连接存在,不代表跨 Broker 进程重启存在;Mosquitto 的 persistence、persistence_location 和自动保存配置决定磁盘恢复,集群产品还会加入复制规则。当前实验配置明确写了 persistence false,只验证同一 Broker 进程内的断线恢复。

离线队列如何回来

实验流程如下:

1
2
3
4
5
6
7
Client ID=mqtt-depth-session, clean=0
→ SUBSCRIBE lab/deep/offline QoS 1
→ 正常断开网络连接,但保留 Session
→ publisher 发布 payload=queued
→ 相同 Client ID 以 clean=0 重连
→ CONNACK Session Present=1
→ 收到 queued

Mosquitto 2.0.18 的实测输出是:

1
session restore: present=1 payload=queued

是否排队、排多少、保存多久还受 QoS、Broker 配置和资源限制影响。MQTT 3.1.1 没有每条消息的标准化到期属性;MQTT 5 才加入 Message Expiry Interval。

Keep Alive 与 Will 的失败时间线

Keep Alive 是客户端声明的最大控制报文发送间隔。只要期间没有别的 MQTT Control Packet,客户端就应发送 PINGREQ;服务器收到后回复 PINGRESP。如果服务器在 Keep Alive 的 1.5 倍时间内没有收到客户端控制报文,必须把网络连接视为失败并断开。MQTT 3.1.1 §3.1.2.10给出了这个 1.5 倍规则。

Will Topic、payload、QoS 和 RETAIN 在 CONNECT 时交给 Broker。I/O 错误、协议错误、Keep Alive 超时等异常结束会触发 Will;客户端发送正常 DISCONNECT 时,Broker 删除 Will 而不发布。

会话恢复、状态边界与 Keep Alive/Will 时间线

本地实验让 Will 客户端声明 Keep Alive=2 秒,收到 CONNACK 后完全静默。最终记录的实际运行在 4.714 秒观察到 payload=timeout。3 秒是服务器判定失败的规范阈值,不是要求实现必须在 3.000 秒精确发布的实时定时器;Mosquitto 的事件循环检查粒度会影响最终观察值。

同一脚本随后建立另一个带 Will 的连接,立即发送 DISCONNECT。观察 2.0 秒没有收到 Will:

1
2
keepalive Will: payload=timeout elapsed=4.714s
graceful DISCONNECT: Will packets=0 during 2.0s observation

完整回环实验

下载以下三个文件:

配置只监听本地回环地址:

1
2
3
4
5
listener 18891 127.0.0.1
allow_anonymous true
persistence false
connection_messages true
log_type all

这份匿名配置只为隔离实验准备,不应复制到对外监听器。启动 Broker:

1
mosquitto -c mqtt-mosquitto.conf -v

在另一个终端运行:

1
python3 mqtt_deep_lab.py

脚本使用 lab/deep/ 主题前缀和 mqtt-depth- Client ID。它先跑不联网的分帧合成测试,再连接本机 Broker。Mosquitto 2.0.18 上的完整输出为:

1
2
3
4
5
6
7
8
9
10
framing: RL(128)=80 01; split=2/15/114 bytes -> 1 packet; coalesced -> 2 packets
mqtt5 fixture: remaining=25 properties=12 total=27 bytes
retained wire: 33 17 00 0f 6c 61 62 2f 64 65 65 70 2f 72 65 74 61 69 6e 00 01 32 33 2e 35
retain scope: established=0 new-subscription=1
session restore: present=1 payload=queued
qos2 duplicate: publish=34 duplicate=3c subscriber-deliveries=1
packet-id reuse after PUBCOMP: payload=after-reuse
keepalive Will: payload=timeout elapsed=4.714s
graceful DISCONNECT: Will packets=0 during 2.0s observation
all checks passed

测试中的 QoS 2 重复不是把应用函数调用两次后比较结果。脚本先发送 DUP=0 的 PUBLISH,收到 PUBREC 后,在 PUBREL 前再次发送相同 PID 和 payload、但把 DUP 置 1,最后才完成 PUBREL/PUBCOMP。订阅端只接受到一个 one-transaction。

Broker 日志也能看到状态顺序:

1
2
3
4
5
6
7
Received PUBLISH ... (d0, q2, r0, m4660, 'lab/deep/qos2', ...)
Sending PUBREC ... (m4660, rc0)
Received PUBLISH ... (d1, q2, r0, m4660, 'lab/deep/qos2', ...)
Sending PUBREC ... (m4660, rc0)
Received PUBREL ... (Mid: 4660)
Sending PUBLISH ... qos2-sub (d0, q2, r0, m1, ...)
Sending PUBCOMP ... (m4660)

日志中的上游十进制 m4660 就是 0x1234;下游使用 m1。这同时验证了重复处理和两个方向独立分配 Packet Identifier。

MQTT 5.0:先读 Properties Length

MQTT 5 保留了固定头和 Remaining Length,在多种控制报文的 Variable Header 中加入 Properties。Properties 区域自身先用一个 Variable Byte Integer 表示长度;解析器必须先限定整个报文,再限定 Properties 子区域,最后按报文类型验证允许出现的属性、数据类型和重复次数。

一条精确的 MQTT 5 PUBLISH

下面是合成并经过脚本长度断言的 QoS 1 PUBLISH:Topic=lab/v5,PID=0x0010,Message Expiry Interval=60,Content Type=text,payload=ok。

1
2
3
4
5
6
7
32 19                         # PUBLISH q1;Remaining Length=25
00 06 6c 61 62 2f 76 35 # Topic Name = lab/v5
00 10 # Packet Identifier
0c # Properties Length = 12
02 00 00 00 3c # Message Expiry Interval = 60
03 00 04 74 65 78 74 # Content Type = "text"
6f 6b # payload = "ok"

长度核算:

1
2
3
Topic field 8 + PID 2 + Properties-Length field 1
+ Properties 12 + payload 2 = Remaining Length 25
fixed header 2 + 25 = total 27 bytes

这段 fixture 没有发送给 Broker;它只验证编码和长度。MQTT 5 Properties 的类型、出现位置和重复限制以 MQTT 5.0 OASIS 标准各报文小节为准。

Clean Start 与 Session Expiry Interval

MQTT 5 把 MQTT 3.1.1 的 Clean Session 拆成两个维度:

Clean Start Session Expiry 结果
1 0 丢弃旧 Session;新 Session 在连接结束时到期
1 大于 0 丢弃旧 Session;创建的新 Session 可在断开后保留
0 0 尝试恢复旧 Session;本次断开后立即到期
0 大于 0 尝试恢复旧 Session,并按间隔继续保留

Clean Start=0 不保证一定存在旧 Session;CONNACK 的 Session Present 仍然是判断依据。Session Expiry Interval 未出现在 CONNECT 时默认为 0。客户端还可以在 DISCONNECT 中调整到期时间,但不能把它从 CONNECT 时的 0 增大为非零。MQTT 5.0 §3.1.2.11.2定义了这些组合。

Receive Maximum 控制在途 QoS 1/2 PUBLISH

CONNECT 与 CONNACK 都可以携带 Receive Maximum。发送端据此建立额度,每发一个 QoS 1/2 PUBLISH 消耗一份;收到 PUBACK 或 PUBCOMP 后归还,收到失败原因码不小于 0x80 的 PUBREC 时也归还。成功 PUBREC 不会归还 QoS 2 的额度,所以进入 PUBREL/PUBCOMP 阶段仍占用窗口。值 0 是 Protocol Error,QoS 0 不消耗这个额度。

若客户端在 CONNECT 声明 Receive Maximum=10,Broker 最多让 10 个下行 QoS 1/2 PUBLISH 同时处于未完成状态。它不是离线队列上限,也不是每秒消息数限制。把 Receive Maximum 设成 1 可以串行化在途确认,却不会自动让应用处理也变成串行事务。

还要区分两个方向:客户端声明的是自己愿意接收多少,服务器在 CONNACK 声明的是服务器愿意从客户端并发接收多少。额度在每次新网络连接时重新初始化,不属于跨连接保存的 Session 状态。相关属性与归还规则见 MQTT 5.0 §3.1.2.11.4、§3.2.2.3.2 与 §4.9。

Subscription Options 的一个字节

MQTT 5 每个订阅 Filter 后有 Subscription Options:

位 含义
1–0 Maximum QoS
2 No Local
3 Retain As Published
5–4 Retain Handling
7–6 保留,必须为 0

例如 0x2d = 0010 1101b 表示 Maximum QoS=1、No Local=1、Retain As Published=1、Retain Handling=2(建立订阅时不发送 retained message)。No Local 不能用于 Shared Subscription。

Retain As Published=0 时,正常转发把 RETAIN 清零;值为 1 时保留发布时的 RETAIN。无论该位如何,因建立普通订阅而发送的 retained message 都带 RETAIN=1。Retain Handling 决定建立或替换订阅时是否补发保留消息。MQTT 5.0 §3.8.3.1列出了位级定义。

Shared Subscription 的标准语义

MQTT 5 标准共享订阅的 Filter 格式是:

1
$share/{ShareName}/{filter}

例如两个 worker 都订阅:

1
$share/image-workers/jobs/image/+

每条匹配发布只选择组内一个 Session,而普通非共享订阅会给每个匹配 Session 各发一份。服务器可以自行选择组内成员,跨多个消费者不承诺全局顺序。共享订阅建立时不发送 retained messages;组中最后一个 Session 离开后,共享订阅及其未交付消息结束。MQTT 5.0 §4.8.2给出了这些约束。

MQTT 3.1.1 标准没有 $share/...。一些 Broker 在 3.1.1 连接上提供同名扩展,部署时应查对应产品与版本文档,不能从抓到 $share 就推断协议协商级别一定是 5。

Message Expiry 与 Will Delay

Message Expiry Interval 让发布者声明消息剩余寿命。Broker 转发时要把剩余秒数更新,过期消息不再交付。它很适合离线队列中的时效命令,例如“十秒内开门”,但应用仍应校验自己的 expires_at,因为跨系统桥接和旧协议节点可能改变路径。

Will Delay Interval 允许 Broker 在网络断开后等待一段时间再发布 Will,使短暂重连不必马上广播离线。若 Session Expiry 先到期,Will 会在会话结束时发布,因此配置时要一起考虑两个间隔。该属性属于 MQTT 5;本文 3.1.1 回环实验没有 Will Delay。

ACL:主题树就是授权边界

下面是 Mosquitto acl_file 的一个精确例子。假设设备用户名就是设备编号 42:

1
2
3
4
5
6
7
8
9
10
pattern write devices/%u/telemetry/#
pattern write devices/%u/status/#
pattern read devices/%u/commands/#
pattern read devices/%u/config

user backend
topic read devices/+/telemetry/#
topic read devices/+/status/#
topic write devices/+/commands/#
topic write devices/+/config

%u 必须独占一个主题层级;%c 可绑定 Client ID。user backend 后的规则按认证用户名匹配,不是按 Client ID。Mosquitto 的 deny 规则会先于允许规则处理,完整语法以当前 mosquitto.conf 官方手册为准。

对应 Broker 配置至少包括:

1
2
3
4
listener 8883
allow_anonymous false
password_file /etc/mosquitto/passwd
acl_file /etc/mosquitto/acl

生产环境还要配置 TLS 服务器证书、客户端的服务器身份校验和凭据轮换。TLS 解决传输机密性与端点认证的一部分;哪个身份能读写哪个 Topic,仍由 ACL 或授权插件决定。

安全测试时可以把检查拆成六个独立问题:

  1. 未认证连接是否成功;
  2. 已认证身份能否订阅自己的主题;
  3. 通配符是否意外扩大到其他设备;
  4. 能否向其他设备的命令主题发布;
  5. 是否能读取历史 retained 数据或恢复旧 Session;
  6. payload 中伪造 device_id 能否绕过应用层对象授权。

一个 CONNECT 成功只说明连接认证通过。SUBACK、实际转发、PUBLISH 确认和业务动作分别需要自己的证据。

应用去重:把副作用放进事务边界

控制消息可以采用稳定的业务信封:

1
2
3
4
5
6
7
8
9
{
"schema": 3,
"producer": "control-plane-a",
"request_id": "01J8...",
"issued_at": "2026-09-22T10:00:00Z",
"expires_at": "2026-09-22T10:00:30Z",
"operation": "unlock",
"resource": "door-42"
}

消费者数据库建立唯一键:

1
2
3
4
5
6
7
8
9
10
create table command_inbox (
tenant_id text not null,
producer text not null,
request_id text not null,
payload_hash bytea not null,
state text not null,
result jsonb,
created_at timestamptz not null,
primary key (tenant_id, producer, request_id)
);

处理流程应让去重记录和数据库副作用处于同一事务:

1
2
3
4
5
6
7
8
BEGIN
→ INSERT inbox;唯一键冲突时读取已有结果
→ 校验 resource 是否属于已认证租户
→ 检查 expires_at 与 schema
→ 执行业务状态变更
→ 写 outbox/result
→ 标记 inbox completed
COMMIT

若副作用是无法纳入数据库事务的物理动作,可以让设备状态机按 request_id 保存执行记录,并把“已接受”“执行中”“已完成”设计成可恢复状态。简单地先开锁、再写去重表,会在两步之间崩溃时留下重复窗口。

去重记录保留时间至少覆盖系统可能重放消息的最长窗口:Session Expiry、离线队列保存时间、桥接积压、生产者重试和业务补偿周期。MQTT 5 Message Expiry 可以缩短传输层窗口,但不能替代消费者的业务规则。

Retained Message 通常适合“当前状态”,不适合具有一次性副作用的命令。若业务确实通过 retained 命令驱动设备,上线后的首次订阅会重新拿到它;应用必须检查 request_id 和过期时间,并在协议设计中明确由谁清除 retained 数据。

抓包与故障定位顺序

Wireshark 可以使用 mqtt 显示过滤器。非默认端口未自动识别时,通过 Decode As 指定 MQTT。排查一条“重复消息”时,我会按这个顺序记录:

1
2
3
4
5
6
7
8
协议版本与连接五元组
→ Client ID、Clean Session/Clean Start、Session Present
→ PUBLISH 的方向、Topic、QoS、DUP、RETAIN、Packet Identifier
→ 对应 ACK 是否在断线前出现
→ 重连后 Session 是否恢复
→ Broker 是否为下游分配了不同 Packet Identifier
→ payload 中是否存在稳定 request_id
→ 应用去重表和副作用提交时间

几个常见现象可以快速映射到协议层:

现象 优先检查
连接后立刻挤掉另一台设备 Client ID 冲突
重连成功但不再收到消息 Session Present=0 后没有重订阅
新订阅立刻收到旧值 retained message 与 RETAIN=1
QoS 1 偶发重复 PUBACK 丢失、Session 恢复、应用主动重发
离线消息在 Broker 重启后消失 persistence 配置与磁盘保存
Keep Alive 设置很小却延迟离线 Broker 检查粒度、网络路径、Will Delay(MQTT 5)
两个 shared worker 都处理同一业务请求 重叠共享组、同时存在普通订阅、应用层重发

结语

MQTT 的主题模型确实简单,可靠性语义则建立在几组严格的局部状态上:TCP 流需要增量分帧;Packet Identifier 只属于当前交换;QoS 分别运行在发布者—Broker和 Broker—订阅者两段;Session、Broker 磁盘持久化、Retained Message、Will 各自保存不同内容。

把这些边界画清楚后,很多抓包现象就能直接解释。最后一道边界仍在应用:消息传输完成,不等于副作用已经以幂等方式提交。稳定的 request_id、授权绑定、过期策略和事务内 inbox/outbox,才是控制类 MQTT 系统真正的可靠性基础。

参考资料


MQTT协议
https://g1at.github.io/2023/08/09/MQTT/
作者
g0at
发布于
2023年8月9日
许可协议