MQTT 协议详解
MQTT 协议详解
MQTT(Message Queuing Telemetry Transport)—— 物联网时代最主流的轻量级消息传输协议
目录
- MQTT 概述
- MQTT 协议核心概念
- MQTT 报文格式详解
- 各类型报文结构详解
- MQTT 5.0 新特性
- MQTT 服务质量(QoS)机制
- MQTT 会话与持久化
- MQTT 安全机制
- MQTT 应用场景
- MQTT 示例代码
- MQTT 与其他协议对比
- MQTT 最佳实践
- 附录:MQTT 速查表
1. MQTT 概述
1.1 什么是 MQTT
MQTT(Message Queuing Telemetry Transport,消息队列遥测传输) 是一种基于 发布/订阅(Pub/Sub) 模式的轻量级消息传输协议。它由 IBM 的 Andy Stanford-Clark 和 Arcom(现为 Eurotech)的 Arlen Nipper 在 1999 年 共同发明,最初用于石油管道遥测系统——需要在极低带宽和高延迟的不稳定网络环境下传输数据。
1.2 发展历程
1999 ── MQTT v1.0 IBM 内部使用,专为石油管道监控设计
│
2010 ── MQTT v3.1 IBM 将 MQTT 免费授权,开始推广
│
2014 ── MQTT v3.1.1 成为 OASIS 标准(国际标准)
│ 轻量级物联网协议的事实标准
│
2016 ── ISO 标准 ISO/IEC 20922 标准化
│
2019 ── MQTT v5.0 ⭐ 重大版本更新
│ 新增会话过期、原因码、用户属性、增强认证等
│
2024+ ─ 持续演进 OASIS MQTT 技术委员会持续维护
1.3 核心特点
| 特点 | 说明 |
|---|---|
| 轻量级 | 固定报头最小仅 2 字节,协议开销极低 |
| 发布/订阅 | 解耦发送方和接收方,实现一对多通信 |
| 三种 QoS | 0(至多一次)、1(至少一次)、2(恰好一次) |
| 双向通信 | 客户端和服务端均可主动发送消息 |
| 遗嘱消息 | 异常断连时自动发布预定义消息 |
| 保留消息 | 新订阅者立即获取最新状态 |
| 会话持久化 | 断线重连后恢复未传递的消息 |
| 最低功耗 | 适合电池供电的嵌入式设备 |
| TLS/SSL 加密 | 支持传输层安全加密 |
1.4 一句话总结
MQTT 就像物联网世界的"微信群"——设备加入群聊(订阅主题),在群里发消息(发布消息),群管理员(Broker)负责转发消息给所有需要的人。而且这个微信群还支持"重发历史消息"(保留消息)、“别人断线了通知你”(遗嘱消息)等高级功能。
2. MQTT 协议核心概念
2.1 发布/订阅模型(Pub/Sub)
MQTT 与传统的客户端/服务器(C/S)请求-响应模型不同,采用发布/订阅模式:
┌─────────────┐
│ MQTT Broker │
│ (消息代理) │
└──────┬──────┘
┌───────────────┼───────────────┐
│ │ │
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│ 客户端 A │ │ 客户端 B │ │ 客户端 C │
│ 发布者 │ │ 订阅者 │ │ 订阅+发布 │
└──────────┘ └──────────┘ └──────────┘
发布 topic/temp 订阅 topic/temp 同时做两件事
核心优势:
| 特性 | 说明 |
|---|---|
| 空间解耦 | 发布者和订阅者不需要知道对方的存在 |
| 时间解耦 | 发布者和订阅者不需要同时在线 |
| 一对多 | 一个发布者可以被多个订阅者接收 |
2.2 主题(Topic)
主题是 MQTT 中消息路由的关键。发布者发布消息到某个主题,订阅者订阅感兴趣的主题。
主题的层级结构:
home/livingroom/temperature
↑ ↑ ↑
第一级 第二级 第三级
通配符:
| 通配符 | 说明 | 示例 |
|---|---|---|
+(单级通配符) |
匹配一个任意层级的主题 | home/+/temperature 匹配 home/livingroom/temperature、home/bedroom/temperature 等 |
#(多级通配符) |
匹配剩余所有层级 | home/# 匹配 home/livingroom/temperature、home/bedroom/light 等 |
通配符使用规则:
✅ 正确用法:
sensor/+/temperature 单级通配符
home/# 多级通配符(必须在末尾)
/+/test 以斜杠开头的单级通配符
❌ 错误用法:
sensor/#/temperature 多级通配符不在末尾
home/#/test 多级通配符后还有层级
2.3 服务质量(QoS)
MQTT 定义了三种服务质量等级(详见第 6 章):
| QoS | 等级 | 可靠性 | 应用场景 |
|---|---|---|---|
| QoS 0 | 至多一次 | 最低,发送即忘 | 传感器数据上报(温度、湿度) |
| QoS 1 | 至少一次 | 中等,保证到达但可能重复 | 控制命令(开关灯) |
| QoS 2 | 恰好一次 | 最高,保证不重复不丢失 | 关键数据(支付消息、告警) |
2.4 MQTT Broker
Broker 是 MQTT 系统的核心服务器,负责消息的路由和转发。
| 主流 Broker | 语言 | 特点 |
|---|---|---|
| EMQX | Erlang | 高并发、可集群、大规模部署首选 |
| Mosquitto | C | 轻量级、标准实现,适合嵌入式网关 |
| NanoMQ | C | 超轻量、边缘计算场景 |
| VerneMQ | Erlang | 高可用、多租户 |
| HiveMQ | Java | 企业级、商业支持 |
| RabbitMQ | Erlang | 通用消息队列,支持 MQTT 插件 |
3. MQTT 报文格式详解
3.1 通用报文结构
所有 MQTT 控制报文都由 固定报头(Fixed Header) 和 可变报头(Variable Header) 以及 有效载荷(Payload) 三部分组成,其中可变报头和有效载荷部分可选。
┌──────────────────────────────────────────┐
│ 固定报头 (Fixed Header) │ ← 所有报文必须
├──────────────────────────────────────────┤
│ 可变报头 (Variable Header) │ ← 部分报文有
├──────────────────────────────────────────┤
│ 有效载荷 (Payload) │ ← 部分报文有
└──────────────────────────────────────────┘
3.2 固定报头结构(Fixed Header)
固定报头是所有 MQTT 报文的第一部分,结构如下:
Byte 1:
┌───────────┬───────────┬───────────┬───────────┐
│ Bit 7 │ Bit 6 │ Bit 5 │ Bit 4 │
├───────────┼───────────┼───────────┼───────────┤
│ 报文类型 (Packet Type) │
│ (4-bit 无符号整数, 1 ~ 15) │
├───────────┼───────────┼───────────┼───────────┤
│ Bit 3 │ Bit 2 │ Bit 1 │ Bit 0 │
├───────────┼───────────┼───────────┼───────────┤
│ 标志位 (Flags) — 取决于报文类型 │
│ 通常为固定值,仅 PUBLISH 报文的标志位有意义 │
└───────────┴───────────┴───────────┴───────────┘
Byte 2...n:
┌──────────────────────────────────────────┐
│ 剩余长度 (Remaining Length) │
│ 可变字节整数 (Variable Byte Integer) │
│ = 可变报头长度 + 有效载荷长度 │
└──────────────────────────────────────────┘
报文类型(前 4 位)
MQTT 5.0 定义了 15 种 控制报文类型:
| 类型 | 名称 | 方向 | 说明 |
|---|---|---|---|
| 0 | Reserved | — | 保留 |
| 1 | CONNECT | C→S | 客户端请求连接服务端 |
| 2 | CONNACK | S→C | 连接确认 |
| 3 | PUBLISH | C↔S | 发布消息 |
| 4 | PUBACK | C↔S | QoS 1 发布确认 |
| 5 | PUBREC | C↔S | QoS 2 第一步(收到) |
| 6 | PUBREL | C↔S | QoS 2 第二步(释放) |
| 7 | PUBCOMP | C↔S | QoS 2 第三步(完成) |
| 8 | SUBSCRIBE | C→S | 客户端订阅主题 |
| 9 | SUBACK | S→C | 订阅确认 |
| 10 | UNSUBSCRIBE | C→S | 客户端取消订阅 |
| 11 | UNSUBACK | S→C | 取消订阅确认 |
| 12 | PINGREQ | C→S | 心跳请求 |
| 13 | PINGRESP | S→C | 心跳响应 |
| 14 | DISCONNECT | C↔S | 断开连接 |
| 15 | AUTH | C↔S | 认证交换(MQTT 5.0 新增) |
标志位(后 4 位)
只有 PUBLISH 报文(类型 3)的标志位有意义:
PUBLISH 标志位:
┌───────┬───────┬───────┬───────┐
│ DUP │ QoS │ QoS │ RETAIN│
│ (bit3)│ (bit2)│ (bit1)│ (bit0)│
└───────┴───────┴───────┴───────┘
| 标志位 | 含义 |
|---|---|
| DUP | 重复标志。1 表示这是一个重发报文(QoS > 0 时使用) |
| QoS (2 bit) | 服务质量等级:0=至多一次,1=至少一次,2=恰好一次 |
| RETAIN | 保留标志。1 表示该消息会被 Broker 保留,新订阅者立即可收到 |
其他报文类型的标志位必须为固定值,通常是 0。
3.3 可变字节整数(Variable Byte Integer, VBI)
这是 MQTT 协议中非常重要的编码方式,用于表示长度可变的字段(如剩余长度、属性长度等)。
编码规则:
┌─────────────────────────────────────┐
│ Byte 1 │ 7-bit 数据 │ 续位=1 │ ← 高位置1表示还有后续
├─────────────────────────────────────┤
│ Byte 2 │ 7-bit 数据 │ 续位=0 │ ← 高位为0表示最后一个字节
└─────────────────────────────────────┘
- 每字节使用 低 7 位 编码数据
- 每字节 最高位(第 8 位) 表示是否还有后续字节:1=还有,0=结束
- 最大 4 字节,可表示范围:0 ~ 268,435,455(256 MB)
示例:
| 数值 | 编码(十六进制) | 说明 |
|---|---|---|
| 0 | 0x00 |
1 字节 |
| 64 | 0x40 |
1 字节 |
| 127 | 0x7F |
1 字节(最大值) |
| 128 | 0x80, 0x01 |
2 字节 |
| 16,383 | 0xFF, 0x7F |
2 字节 |
| 16,384 | 0x80, 0x80, 0x01 |
3 字节 |
3.4 剩余长度(Remaining Length)
固定报头中的剩余长度字段使用 VBI 编码,表示可变报头长度 + 有效载荷长度。这个字段非常重要,因为接收方需要知道要读取多少字节才能构成一个完整的报文。
3.5 可变报头(Variable Header)
可变报头的内容严格依赖于报文类型,不同报文的字段各不相同。
MQTT 5.0 新增——属性(Properties):
所有 MQTT 5.0 的可变报头末尾都可以携带属性列表:
属性部分结构:
┌──────────────────────────────────────────┐
│ 属性长度 (Properties Length) │
│ └ 可变字节整数,表示后面所有属性的总长度 │
├──────────────────────────────────────────┤
│ 属性 1 │ 属性标识符 │ 属性值 │
│ 属性 2 │ 属性标识符 │ 属性值 │
│ ... │
└──────────────────────────────────────────┘
通用属性类型:
| 属性标识符 | 名称 | 类型 | 适用报文 |
|---|---|---|---|
| 0x01 | Payload Format Indicator | 单字节 | PUBLISH |
| 0x02 | Message Expiry Interval | 4 字节整数 | PUBLISH |
| 0x03 | Content Type | UTF-8 字符串 | PUBLISH |
| 0x08 | Response Topic | UTF-8 字符串 | PUBLISH |
| 0x09 | Correlation Data | 二进制数据 | PUBLISH |
| 0x0B | Subscription Identifier | 可变字节整数 | PUBLISH / SUBSCRIBE |
| 0x11 | Session Expiry Interval | 4 字节整数 | CONNECT / CONNACK / DISCONNECT |
| 0x12 | Assigned Client Identifier | UTF-8 字符串 | CONNACK |
| 0x13 | Server Keep Alive | 2 字节整数 | CONNACK |
| 0x15 | Authentication Method | UTF-8 字符串 | CONNECT / CONNACK / AUTH |
| 0x16 | Authentication Data | 二进制数据 | CONNECT / CONNACK / AUTH |
| 0x17 | Request Problem Information | 单字节 | CONNECT |
| 0x19 | Request Response Information | 单字节 | CONNECT |
| 0x1A | Response Information | UTF-8 字符串 | CONNACK |
| 0x1C | Server Reference | UTF-8 字符串 | CONNACK |
| 0x1F | Reason String | UTF-8 字符串 | 所有确认报文 |
| 0x21 | Receive Maximum | 2 字节整数 | CONNECT / CONNACK |
| 0x22 | Topic Alias Maximum | 2 字节整数 | CONNECT / CONNACK |
| 0x23 | Topic Alias | 2 字节整数 | PUBLISH |
| 0x24 | Maximum QoS | 单字节 | CONNACK |
| 0x25 | Retain Available | 单字节 | CONNACK |
| 0x26 | User Property | UTF-8 字符串对 | 所有报文(可重复) |
| 0x27 | Maximum Packet Size | 4 字节整数 | CONNECT / CONNACK |
| 0x28 | Wildcard Subscription Available | 单字节 | CONNACK |
| 0x29 | Subscription Identifier Available | 单字节 | CONNACK |
| 0x2A | Shared Subscription Available | 单字节 | CONNACK |
4. 各类型报文结构详解
4.1 CONNECT 报文(客户端 → 服务端)
CONNECT 报文是客户端与服务端建立连接时发送的第一个报文。
固定报头:
0x10 + 剩余长度
可变报头:
┌──────────────────────────────────────────┐
│ 协议名称 (Protocol Name) │
│ 格式: UTF-8 字符串 "MQTT" │
│ 二进制: 0x00 0x04 'M' 'Q' 'T' 'T' │ ← 6 字节
├──────────────────────────────────────────┤
│ 协议版本 (Protocol Version) │
│ MQTT 5.0 = 5, MQTT 3.1.1 = 4 │ ← 1 字节
├──────────────────────────────────────────┤
│ 连接标志 (Connect Flags) │
│ ┌─────┬─────┬─────┬─────┬─────┬─────┬────┐│
│ │User │Pass │Will │Will │Will │Clean│Resv││
│ │Name │word │Retain│QoS │Flag │Start│=0 ││
│ │(1) │(1) │(1) │(2) │(1) │(1) │(1) ││
│ └──7──┴──6──┴──5──┴4─3─┴──2──┴──1──┴──0─┘│ ← 1 字节
├──────────────────────────────────────────┤
│ 保持连接 (Keep Alive) │
│ 单位: 秒, 0=禁用 Keep Alive 检测 │ ← 2 字节
├──────────────────────────────────────────┤
│ 属性 (Properties) │ ← 变长
└──────────────────────────────────────────┘
有效载荷:
┌──────────────────────────────────────────┐
│ 客户端标识符 (Client ID) │ ← UTF-8 字符串
├──────────────────────────────────────────┤
│ 遗嘱属性 (Will Properties) [可选] │ ← Will Flag=1
├──────────────────────────────────────────┤
│ 遗嘱主题 (Will Topic) [可选] │ ← Will Flag=1
├──────────────────────────────────────────┤
│ 遗嘱载荷 (Will Payload) [可选] │ ← Will Flag=1
├──────────────────────────────────────────┤
│ 用户名 (User Name) [可选] │ ← User Name Flag=1
├──────────────────────────────────────────┤
│ 密码 (Password) [可选] │ ← Password Flag=1
└──────────────────────────────────────────┘
MQTT 3.1.1 与 5.0 的区别:MQTT 5.0 使用 Clean Start(而非 Clean Session),并在可变报头中增加了 Properties。
4.2 CONNACK 报文(服务端 → 客户端)
固定报头:
0x20 + 剩余长度
可变报头:
┌──────────────────────────────────────────┐
│ 连接确认标志 (Connect Acknowledge Flags) │
│ Bit 0: Session Present │ ← 1 字节
│ Bit 1~7: 保留为 0 │
├──────────────────────────────────────────┤
│ 原因码 (Reason Code) │
│ 0x00 = 成功 │ ← 1 字节
│ 0x84 = 不支持的协议版本 │
│ 0x85 = 客户端标识符无效 │
│ 0x86 = 用户名或密码错误 │
│ 更多原因码见附录 │
├──────────────────────────────────────────┤
│ 属性 (Properties) │ ← 变长
└──────────────────────────────────────────┘
4.3 PUBLISH 报文(发布消息)
固定报头:
Byte 1: 0b0011_xxxx
│ │└─ RETAIN
│ └── QoS
└────────── DUP
Byte 2+: 剩余长度
可变报头:
┌──────────────────────────────────────────┐
│ 主题名 (Topic Name) │ ← UTF-8 字符串
│ 格式: 2字节长度 + UTF-8编码的主题字符串 │
├──────────────────────────────────────────┤
│ 报文标识符 (Packet Identifier) [可选] │ ← QoS > 0 时存在
│ │ ← 2 字节
├──────────────────────────────────────────┤
│ 属性 (Properties) │ ← 变长
└──────────────────────────────────────────┘
有效载荷:
┌──────────────────────────────────────────┐
│ 应用消息 (Application Message) │
│ 二进制数据, 内容和格式由应用自行定义 │
│ 长度 = 剩余长度 - 可变报头长度 │
└──────────────────────────────────────────┘
报文标识符的用法:
PUBLISH (QoS 0): 无报文标识符
PUBLISH (QoS 1): → PUBACK (Packet Identifier 必须相同)
PUBLISH (QoS 2): → PUBREC → PUBREL → PUBCOMP (所有报文使用相同 Packet Identifier)
4.4 PUBACK 报文(QoS 1 确认)
固定报头:
0x40 + 剩余长度
可变报头(MQTT 5.0):
┌──────────────────────────────────────────┐
│ 报文标识符 (Packet Identifier) │ ← 2 字节
├──────────────────────────────────────────┤
│ 原因码 (Reason Code) │ ← 1 字节
│ 0x00 = 成功 │
│ 0x10 = 无匹配订阅者 │
├──────────────────────────────────────────┤
│ 属性 (Properties) │ ← 可选,变长
└──────────────────────────────────────────┘
4.5 SUBSCRIBE 报文(订阅主题)
固定报头:
0x82 + 剩余长度
可变报头:
┌──────────────────────────────────────────┐
│ 报文标识符 (Packet Identifier) │ ← 2 字节
├──────────────────────────────────────────┤
│ 属性 (Properties) │ ← 变长
└──────────────────────────────────────────┘
有效载荷:
┌──────────────────────────────────────────┐
│ 主题过滤器 1 │
│ ├── 主题长度 (2 字节) │
│ ├── 主题字符串 │
│ └── 订阅选项 (1 字节) │
│ ├── Bit 0-1: QoS │
│ ├── Bit 2: No Local │ (MQTT 5.0)
│ ├── Bit 3: Retain As Published │ (MQTT 5.0)
│ └── Bit 4-5: Retain Handling │ (MQTT 5.0)
├──────────────────────────────────────────┤
│ 主题过滤器 2 │
│ ... │
└──────────────────────────────────────────┘
MQTT 5.0 新增订阅选项详解:
| 选项 | 位域 | 说明 |
|---|---|---|
| No Local | Bit 2 | 1=自己发布的消息不会回传给自己的订阅 |
| Retain As Published | Bit 3 | 保留消息是否保持原始格式转发 |
| Retain Handling | Bit 4-5 | 0=任何时候都发送保留消息;1=新订阅时才发送;2=不发送保留消息 |
4.6 SUBACK 报文(订阅确认)
固定报头:
0x90 + 剩余长度
可变报头:
┌──────────────────────────────────────────┐
│ 报文标识符 (Packet Identifier) │ ← 2 字节
├──────────────────────────────────────────┤
│ 属性 (Properties) │ ← 变长
└──────────────────────────────────────────┘
有效载荷:
┌──────────────────────────────────────────┐
│ 原因码列表 (每个主题过滤器一个) │
│ 0x00 = QoS 0 - 订阅成功,最高可接受 QoS 0│
│ 0x01 = QoS 1 - 订阅成功,最高可接受 QoS 1│
│ 0x02 = QoS 2 - 订阅成功,最高可接受 QoS 2│
│ 0x80 = 订阅失败/未授权 │
└──────────────────────────────────────────┘
4.7 PINGREQ / PINGRESP 报文(心跳)
PINGREQ 报文(客户端 → 服务端):
0xC0 0x00 ← 只有固定报头,没有可变报头和有效载荷
PINGRESP 报文(服务端 → 客户端):
0xD0 0x00 ← 只有固定报头
- PINGREQ/PINGRESP 是维持连接活性最基本的机制
- 客户端在 Keep Alive 时间内未发送任何报文时,必须发送 PINGREQ
- 服务端收到 PINGREQ 后应立即回复 PINGRESP
4.8 DISCONNECT 报文
固定报头:
0xE0 + 剩余长度
可变报头(MQTT 5.0):
┌──────────────────────────────────────────┐
│ 原因码 (Reason Code) │ ← 1 字节 (可选)
│ 0x00 = 正常断开 │
│ 0x04 = 以 Will 消息的形式断开 │
├──────────────────────────────────────────┤
│ 属性 (Properties) │ ← 变长 (可选)
│ 可以携带 Session Expiry Interval │
│ 可以携带 Reason String │
│ 可以携带 User Property │
└──────────────────────────────────────────┘
4.9 AUTH 报文(MQTT 5.0 新增)
用于增强认证交换,支持 SCRAM、Kerberos 等更安全的认证方式。
固定报头:
0xF0 + 剩余长度
可变报头:
┌──────────────────────────────────────────┤
│ 原因码 (Reason Code) │
│ 0x00 = 成功 │
│ 0x18 = 继续认证 │
│ 0x19 = 重新认证 │
├──────────────────────────────────────────┤
│ 属性 (Properties) │
│ 包含 Authentication Method │
│ 包含 Authentication Data │
└──────────────────────────────────────────┘
5. MQTT 5.0 新特性
MQTT 5.0 是自 2014 年 3.1.1 版本以来的首次大版本更新,引入了大量新特性:
| 新特性 | 说明 | 解决了什么问题 |
|---|---|---|
| 原因码(Reason Code) | 所有确认报文(CONNACK、SUBACK 等)都携带原因码 | MQTT 3.1.1 只有成功/失败两种结果,无法表达"失败原因" |
| 属性(Properties) | 所有报文类型都可携带键值对属性 | 大大扩展了协议的扩展能力 |
| 会话过期(Session Expiry) | 会话可以设置过期时间,不再永久保存 | 服务端存储管理更高效 |
| 报文过期(Message Expiry) | 每个消息可以设置过期时间 | 避免过期消息继续传递 |
| 主题别名(Topic Alias) | 用短整型替代长主题字符串 | 减少长主题场景的开销 |
| 用户属性(User Property) | 支持自定义键值对元数据 | 应用层可附加自定义信息 |
| 订阅选项(Subscription Options) | No Local、Retain As Published、Retain Handling | 订阅行为更灵活可控 |
| 增强认证(Enhanced Authentication) | 支持 SCRAM、Kerberos 等多轮认证 | 安全能力大幅提升 |
| 请求/响应模式 | Response Topic + Correlation Data | 支持 RPC 风格通信 |
| 服务端断开(Server-initiated DISCONNECT) | 服务端可以主动断开连接 | 服务端管理能力增强 |
| 流量控制(Flow Control) | Receive Maximum 控制并行消息数 | 避免消息过载 |
| 最大报文大小协商 | 客户端和服务端可协商最大报文尺寸 | 适配不同设备能力 |
6. MQTT 服务质量(QoS)机制
6.1 QoS 0 —— 至多一次(At Most Once)
发送者 ─────PUBLISH─────► 接收者
(发送即忘,不确认)
- 最快速、最轻量
- 消息可能丢失(网络中断时)
- 不需要存储报文
- 典型场景:传感器每秒上报温度,偶尔丢一帧没关系
6.2 QoS 1 —— 至少一次(At Least Once)
发送者 ─────PUBLISH─────► 接收者
◄─────PUBACK───── 发送者
(收到确认前可能重发)
- 保证消息到达至少一次
- 消息可能重复(重发导致)
- 发送者需要存储报文直到收到 PUBACK
- 典型场景:控制命令(开关灯),重复执行一次不影响
流程详解:
1. 发送者发送 PUBLISH,启动定时器,等待 PUBACK
2. 接收者收到 PUBLISH,发送 PUBACK
3a. 发送者收到 PUBACK → 报文发送完成
3b. 超时未收到 PUBACK → 发送者重发 PUBLISH(DUP=1)
6.3 QoS 2 —— 恰好一次(Exactly Once)
发送者 ──────PUBLISH──────► 接收者
◄─────PUBREC─────── 收到
──────PUBREL───────► 释放
◄─────PUBCOMP────── 完成
- 最高可靠性,保证消息不重复不丢失
- 需要 4 次握手(PUBLISH → PUBREC → PUBREL → PUBCOMP)
- 发送者和接收者都需要存储报文状态
- 开销最大(网络流量和存储)
- 典型场景:支付交易、关键告警、固件升级指令
流程详解:
发送者: 接收者:
1. PUBLISH (QoS=2) ──────────────► 存储报文,标记"已收到"
2. ◄─────────────────── PUBREC 存储报文标识符
3. 丢弃原 PUBLISH,存储报文标识符
4. PUBREL ───────────────────────► 收到 PUBREL
5. 交付消息给订阅者
6. ◄─────────────────── PUBCOMP 清理报文存储
7. 清理报文标识符存储
6.4 QoS 降级规则
当发布者的 QoS 等级高于订阅者请求的 QoS 等级时,Broker 会降级:
发布者 QoS 2 ──→ Broker ──→ 订阅者请求 QoS 1
└── 实际交付 QoS 1(降级)
规则:最终交付的 QoS = min(发布者 QoS, 订阅者请求的 QoS)
7. MQTT 会话与持久化
7.1 Clean Start 与 Session Expiry(MQTT 5.0)
Clean Start(全新开始):
| Clean Start | 行为 |
|---|---|
| 1 | 服务端丢弃该客户端的所有现有会话数据,开启全新会话 |
| 0 | 服务端尝试恢复客户端之前的会话 |
Session Expiry Interval(会话过期时间):
| 值 | 含义 |
|---|---|
| 0 | 连接断开时会话立即过期(等价于 MQTT 3.1.1 的 Clean Session=1) |
| 0xFFFFFFFF | 会话永不过期 |
| 其他 | 会话在断开后 N 秒过期 |
7.2 会话中保存的数据
服务端为每个连接的客户端维护以下数据:
客户端会话数据:
├── 已订阅的主题列表及订阅选项
├── 所有 QoS 1 和 QoS 2 的未确认消息
├── 所有 QoS 2 的未完成报文标识符
└── 遗嘱消息(如果设定了 Will)
7.3 保留消息(Retained Message)
保留消息是 MQTT 非常有特色的功能:
发布者:
发送 PUBLISH (RETAIN=1) topic: "sensor/temp" payload: "25.5°C"
Broker 行为:
1. 将消息转发给当前所有订阅者
2. 在 Broker 上存储该主题的最后一条保留消息
新订阅者(之后的):
订阅 topic: "sensor/temp"
→ 立即收到保留消息 "25.5°C",无需等待下次发布
→ 然后继续接收后续发布的实时消息
取消保留:
向该主题发送一条空 payload 的保留消息 (payload="", RETAIN=1)
7.4 遗嘱消息(Will Message)
遗嘱消息是 MQTT 异常断连的最后保障:
客户端在 CONNECT 时设定:
Will Topic: "device/status"
Will Payload: "offline"
Will QoS: 1
Will Retain: true
正常场景:
客户端主动发送 DISCONNECT → Broker 不发送遗嘱消息
异常断连场景:
客户端网络断开 / 掉电 / 崩溃
→ Broker 检测到 Keep Alive 超时
→ Broker 自动发布遗嘱消息给所有订阅者
→ "device/status" → "offline"(其他设备立即可知该设备离线)
8. MQTT 安全机制
8.1 传输层安全(TLS/SSL)
| 安全等级 | 说明 | 适用场景 |
|---|---|---|
| 无加密 | 明文传输,性能最好 | 内部网络 / 开发测试 |
| 单向 TLS | 服务端证书验证,客户端不验证 | 大多数生产环境 |
| 双向 TLS | 服务端和客户端证书互相验证 | 金融、医疗等高安全场景 |
8.2 应用层安全
| 机制 | 说明 |
|---|---|
| 用户名/密码 | CONNECT 报文中的 User Name 和 Password 字段 |
| Token 认证 | 使用密码字段传输 JWT Token |
| 增强认证 | MQTT 5.0 的 AUTH 报文支持 SCRAM、Kerberos 等多轮认证协议 |
8.3 授权与 ACL
Broker 端的访问控制列表(ACL):
# Mosquitto ACL 示例
# 用户 alice 可以发布和订阅 sensor 下的所有主题
user alice
topic read sensor/#
topic write sensor/#
# 用户 bob 只能订阅
user bob
topic read home/#
# 匿名用户只能订阅 public 主题
topic read public/#
9. MQTT 应用场景
9.1 物联网(IoT)—— 最大应用场景
┌─────────┐ ┌─────────┐ ┌─────────┐
│温度传感器 │ │ 烟雾传感器│ │ 门锁 │
│(发布) │ │ (发布) │ │ (订阅) │
└────┬─────┘ └────┬─────┘ └────┬─────┘
│ │ │
▼ ▼ ▼
┌──────────────────────────────────────┐
│ MQTT Broker │
└──────┬────────────────────┬──────────┘
│ │
▼ ▼
┌──────────┐ ┌──────────┐
│ 手机App │ │ 云端服务 │
│ (订阅) │ │ (订阅/存储 │
└──────────┘ └──────────┘
| IoT 领域 | 具体应用 |
|---|---|
| 智能家居 | 灯光控制、温湿度监测、安防报警、智能窗帘 |
| 工业物联网 | 设备状态监控、PLC 数据采集、预测性维护 |
| 车联网 | 车辆定位上报、远程 OTA 升级、故障诊断 |
| 智慧农业 | 土壤湿度监测、自动灌溉控制、气象站数据 |
| 智慧城市 | 路灯控制、垃圾箱满溢检测、停车位管理 |
| 医疗健康 | 可穿戴设备数据、远程病人监护、医疗设备状态 |
9.2 即时通讯 / 消息推送
- 轻量级聊天应用(群聊使用共享订阅)
- 手机推送通知
- 实时消息广播(如新闻推送)
9.3 移动应用
- 手机端的实时数据同步
- 定位追踪上报
- 协作类应用的实时更新
9.4 边缘计算
边缘网关 (Edge Gateway)
├── 从传感器采集数据
├── 本地做数据过滤和聚合
├── MQTT 上报云端
└── 本地执行规则引擎(断网时自动处理)
9.5 服务器监控与告警
| 场景 | 实现方式 |
|---|---|
| 服务器 CPU/内存监控 | Agent 定时 PUBLISH 指标数据 |
| 服务状态告警 | 异常时 PUBLISH 到告警主题 |
| 日志实时收集 | 各服务 PUBLISH 日志条目 |
| 自动化运维 | 订阅命令主题执行远程指令 |
10. MQTT 示例代码
10.1 Python 示例(使用 paho-mqtt)
安装依赖:
pip install paho-mqtt
示例 1:最简单的发布者和订阅者
发布者(publisher.py):
import paho.mqtt.client as mqtt
import time
import json
# Broker 配置
BROKER_HOST = "broker.emqx.io" # 公共测试 Broker
BROKER_PORT = 1883
TOPIC = "test/temperature"
CLIENT_ID = "publisher_demo_001"
def on_connect(client, userdata, flags, reason_code, properties):
"""连接回调"""
print(f"已连接到 Broker,原因码: {reason_code}")
def on_publish(client, userdata, mid, reason_code, properties):
"""发布回调"""
print(f"消息已发布, mid: {mid}")
def on_disconnect(client, userdata, reason_code, properties):
"""断开回调"""
print(f"已断开连接, 原因码: {reason_code}")
# 1. 创建客户端(MQTT 5.0 使用 CallbackAPIVersion.VERSION2)
client = mqtt.Client(
client_id=CLIENT_ID,
callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
protocol=mqtt.MQTTv5
)
# 2. 绑定回调函数
client.on_connect = on_connect
client.on_publish = on_publish
client.on_disconnect = on_disconnect
# 3. 设置遗嘱消息(可选)
will_topic = "test/status"
will_payload = json.dumps({"device": CLIENT_ID, "status": "offline"})
client.will_set(will_topic, will_payload, qos=1, retain=True)
# 4. 连接 Broker
client.connect(BROKER_HOST, BROKER_PORT, keepalive=60)
client.loop_start()
# 5. 循环发布消息
for i in range(10):
payload = json.dumps({
"device": CLIENT_ID,
"temperature": round(25.0 + i * 0.5, 1),
"humidity": round(60.0 - i * 0.3, 1),
"timestamp": int(time.time())
})
result = client.publish(TOPIC, payload, qos=1, retain=False)
print(f"[{i+1}/10] 已发布: {payload}")
time.sleep(2)
# 6. 发布在线状态(保留消息)
status_payload = json.dumps({"device": CLIENT_ID, "status": "online"})
client.publish("test/status", status_payload, qos=1, retain=True)
# 7. 清理并断开
time.sleep(1)
client.loop_stop()
client.disconnect()
订阅者(subscriber.py):
import paho.mqtt.client as mqtt
import json
BROKER_HOST = "broker.emqx.io"
BROKER_PORT = 1883
TOPIC = "test/temperature"
CLIENT_ID = "subscriber_demo_001"
def on_connect(client, userdata, flags, reason_code, properties):
"""连接成功后自动订阅"""
if reason_code == 0:
print(f"连接成功,正在订阅主题: {TOPIC}")
client.subscribe(TOPIC, qos=1)
else:
print(f"连接失败,原因码: {reason_code}")
def on_message(client, userdata, msg):
"""收到消息回调"""
try:
payload = json.loads(msg.payload)
print(f"[收到消息] 主题: {msg.topic}")
print(f" QoS: {msg.qos}")
print(f" 内容: {json.dumps(payload, indent=2, ensure_ascii=False)}")
print(f" 保留: {msg.retain}")
print("-" * 50)
except Exception as e:
print(f"[收到消息] 主题: {msg.topic}")
print(f" 内容: {msg.payload.decode()}")
print("-" * 50)
def on_subscribe(client, userdata, mid, reason_code_list, properties):
"""订阅回调"""
print(f"订阅成功, mid: {mid}, 原因码列表: {reason_code_list}")
def on_disconnect(client, userdata, reason_code, properties):
print(f"已断开连接")
# 创建客户端
client = mqtt.Client(
client_id=CLIENT_ID,
callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
protocol=mqtt.MQTTv5
)
client.on_connect = on_connect
client.on_message = on_message
client.on_subscribe = on_subscribe
client.on_disconnect = on_disconnect
client.connect(BROKER_HOST, BROKER_PORT, keepalive=60)
# 阻塞方式运行(持续等待消息)
print(f"正在监听主题: {TOPIC}")
print(f"按 Ctrl+C 退出")
try:
client.loop_forever()
except KeyboardInterrupt:
print("用户中断")
client.disconnect()
运行方式:
# 终端 1:启动订阅者
python subscriber.py
# 终端 2:启动发布者
python publisher.py
示例 2:MQTT 5.0 完整功能演示
"""
MQTT 5.0 高级特性演示
包含:主题别名、用户属性、会话过期、请求/响应模式
"""
import paho.mqtt.client as mqtt
from paho.mqtt.properties import Properties
from paho.mqtt.packettypes import PacketTypes
import json
import time
import threading
BROKER = "broker.emqx.io"
PORT = 1883
class MQTT5Demo:
def __init__(self, client_id):
self.client = mqtt.Client(
client_id=client_id,
callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
protocol=mqtt.MQTTv5
)
self.client.on_connect = self.on_connect
self.client.on_message = self.on_message
self.client.on_subscribe = self.on_subscribe
# 启用自动重连
self.client.reconnect_delay_set(min_delay=1, max_delay=120)
def on_connect(self, client, userdata, flags, reason_code, properties):
print(f"[{client._client_id.decode()}] 连接成功 (原因码={reason_code})")
def on_subscribe(self, client, userdata, mid, reason_code_list, properties):
print(f"订阅成功: mid={mid}, codes={reason_code_list}")
def on_message(self, client, userdata, msg):
# 解析用户属性
user_properties = {}
if msg.properties:
for k, v in msg.properties.UserProperty:
user_properties[k] = v
print(f"\n[{client._client_id.decode()}] 收到消息:")
print(f" 主题: {msg.topic}")
print(f" QoS: {msg.qos}")
print(f" 载荷: {msg.payload.decode()[:100]}")
if user_properties:
print(f" 用户属性: {user_properties}")
if msg.properties and msg.properties.CorrelationData:
print(f" CorrelationData: {msg.properties.CorrelationData.hex()}")
def connect(self, clean_start=True, session_expiry=0):
"""连接 Broker,并使用 MQTT 5.0 属性"""
props = Properties(PacketTypes.CONNECT)
props.SessionExpiryInterval = session_expiry
# 请求响应信息
props.RequestResponseInformation = 1
# 请求问题信息
props.RequestProblemInformation = 1
# 附加用户属性
props.UserProperty = [("language", "zh-CN"), ("version", "1.0")]
self.client.connect(BROKER, PORT, keepalive=60,
clean_start=clean_start, properties=props)
def publish_with_properties(self, topic, payload, qos=1, retain=False):
"""使用 MQTT 5.0 属性发布消息"""
props = Properties(PacketTypes.PUBLISH)
# 消息过期时间(30秒后过期)
props.MessageExpiryInterval = 30
# 载荷格式标识(1=UTF-8 文本)
props.PayloadFormatIndicator = 1
# 内容类型
props.ContentType = "application/json"
# 用户属性(可以附加自定义元数据)
props.UserProperty = [
("source", "mqtt5_demo"),
("timestamp", str(int(time.time())))
]
info = self.client.publish(topic, payload, qos=qos, retain=retain,
properties=props)
return info
def publish_request(self, topic, payload, response_topic, correlation_data):
"""MQTT 5.0 请求/响应模式"""
props = Properties(PacketTypes.PUBLISH)
props.ResponseTopic = response_topic
props.CorrelationData = correlation_data
props.UserProperty = [("type", "request")]
self.client.publish(topic, payload, qos=1, properties=props)
def start(self):
self.client.loop_start()
def stop(self):
self.client.loop_stop()
self.client.disconnect()
def demo_request_response():
"""请求/响应模式演示"""
requester = MQTT5Demo("requester_demo")
responder = MQTT5Demo("responder_demo")
# 响应者订阅响应主题(实际场景中通常是固定主题)
RESPONSE_TOPIC = "demo/response"
def responder_on_message(client, userdata, msg):
# 解析请求
req = json.loads(msg.payload)
print(f"\n[响应者] 收到请求: {req}")
# 检查是否有 ResponseTopic
if msg.properties and msg.properties.ResponseTopic:
resp_topic = msg.properties.ResponseTopic
corr_data = msg.properties.CorrelationData
# 发送响应
resp_payload = json.dumps({
"status": "ok",
"result": f"处理了: {req.get('command', 'unknown')}",
"timestamp": time.time()
})
props = Properties(PacketTypes.PUBLISH)
props.CorrelationData = corr_data
props.UserProperty = [("type", "response")]
client.publish(resp_topic, resp_payload, qos=1, properties=props)
print(f"[响应者] 已发送响应到 {resp_topic}")
responder.client.on_message = responder_on_message
responder.connect(clean_start=True)
responder.client.subscribe("demo/request")
responder.start()
requester.connect(clean_start=True)
requester.client.subscribe(RESPONSE_TOPIC)
def requester_on_message(client, userdata, msg):
print(f"[请求者] 收到响应: {msg.payload.decode()}")
if msg.properties and msg.properties.UserProperty:
print(f" 属性: {dict(msg.properties.UserProperty)}")
requester.client.on_message = requester_on_message
requester.start()
time.sleep(1)
# 发送请求
import uuid
corr_id = uuid.uuid4().bytes
requester.publish_request(
"demo/request",
json.dumps({"command": "get_status", "device": "sensor01"}),
RESPONSE_TOPIC,
corr_id
)
time.sleep(3)
requester.stop()
responder.stop()
if __name__ == "__main__":
print("=" * 50)
print("MQTT 5.0 演示程序")
print("=" * 50)
print("\n--- 基本发布/订阅演示 ---")
pub = MQTT5Demo("pub_advanced")
sub = MQTT5Demo("sub_advanced")
sub.connect(clean_start=True, session_expiry=60)
sub.client.subscribe("demo/advanced")
sub.start()
pub.connect(clean_start=True)
pub.start()
time.sleep(1)
# 使用高级属性发布消息
pub.publish_with_properties(
"demo/advanced",
json.dumps({"msg": "Hello MQTT 5.0!", "seq": 1, "time": time.time()}),
qos=1
)
time.sleep(2)
print("\n--- 请求/响应模式演示 ---")
demo_request_response()
pub.stop()
sub.stop()
示例 3:嵌入式设备场景(C 语言 paho 客户端)
/*
* 嵌入式 MQTT 客户端示例(基于 Eclipse Paho C 库)
* 模拟温度传感器设备
* gcc -o mqtt_sensor mqtt_sensor.c -lpaho-mqtt3cs
*/
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <time.h>
#include "MQTTClient.h"
#define BROKER_ADDRESS "tcp://broker.emqx.io:1883"
#define CLIENT_ID "embedded_sensor_001"
#define TOPIC_TEMP "sensor/temperature"
#define TOPIC_STATUS "sensor/status"
#define QOS 1
#define TIMEOUT 10000L
// 模拟读取传感器数据
float read_temperature() {
// 模拟一个在 20~30°C 之间波动的温度值
return 25.0f + (rand() % 100) / 10.0f;
}
int message_arrived(void *context, char *topicName, int topicLen,
MQTTClient_message *message) {
printf("[收到命令] 主题: %s, 内容: %.*s\n",
topicName, message->payloadlen, (char*)message->payload);
MQTTClient_freeMessage(&message);
MQTTClient_free(topicName);
return 1;
}
int main(int argc, char* argv[]) {
MQTTClient client;
MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer;
MQTTClient_willOptions will_opts = MQTTClient_willOptions_initializer;
int rc;
// 1. 创建客户端
MQTTClient_create(&client, BROKER_ADDRESS, CLIENT_ID,
MQTTCLIENT_PERSISTENCE_NONE, NULL);
// 2. 设置回调
MQTTClient_setCallbacks(client, NULL, NULL, message_arrived, NULL);
// 3. 配置连接选项
conn_opts.keepAliveInterval = 20;
conn_opts.cleansession = 1;
conn_opts.username = "";
conn_opts.password = "";
// 4. 设置遗嘱消息
will_opts.topicName = TOPIC_STATUS;
will_opts.message = "offline";
will_opts.qos = 1;
will_opts.retained = 1;
conn_opts.will = &will_opts;
// 5. 连接 Broker
rc = MQTTClient_connect(client, &conn_opts);
if (rc != MQTTCLIENT_SUCCESS) {
printf("连接失败! 错误码: %d\n", rc);
MQTTClient_destroy(&client);
return -1;
}
printf("已连接到 Broker: %s\n", BROKER_ADDRESS);
// 6. 订阅控制主题
MQTTClient_subscribe(client, "sensor/control", QOS);
// 7. 发布在线状态
MQTTClient_message pubmsg = MQTTClient_message_initializer;
pubmsg.payload = "online";
pubmsg.payloadlen = 6;
pubmsg.qos = 1;
pubmsg.retained = 1;
MQTTClient_publishMessage(client, TOPIC_STATUS, &pubmsg, NULL);
// 8. 循环发布温度数据
printf("开始上报温度数据 (按 Ctrl+C 退出)...\n");
for (int i = 0; i < 30; i++) {
float temp = read_temperature();
char buf[32];
snprintf(buf, sizeof(buf), "%.1f", temp);
pubmsg.payload = buf;
pubmsg.payloadlen = strlen(buf);
pubmsg.qos = QOS;
pubmsg.retained = 0;
rc = MQTTClient_publishMessage(client, TOPIC_TEMP, &pubmsg, NULL);
printf("[%d] 温度: %s°C (rc=%d)\n", i + 1, buf, rc);
sleep(5);
}
// 9. 断开连接
MQTTClient_disconnect(client, 10000);
MQTTClient_destroy(&client);
printf("连接已关闭\n");
return 0;
}
示例 4:JavaScript / Node.js 示例
// 安装: npm install mqtt
const mqtt = require('mqtt')
const BROKER = 'mqtt://broker.emqx.io:1883'
const CLIENT_ID = 'nodejs_client_' + Math.random().toString(16).substr(2, 8)
// 创建客户端连接
const client = mqtt.connect(BROKER, {
clientId: CLIENT_ID,
clean: true,
connectTimeout: 4000,
reconnectPeriod: 1000, // 自动重连间隔
will: {
topic: 'device/status',
payload: JSON.stringify({ clientId: CLIENT_ID, status: 'offline' }),
qos: 1,
retain: true
}
})
client.on('connect', () => {
console.log(`已连接到 Broker (clientId: ${CLIENT_ID})`)
// 订阅多个主题
client.subscribe('device/+/status', { qos: 1 }, (err) => {
if (!err) console.log('已订阅 device/+/status')
})
client.subscribe('sensor/#', { qos: 0 })
// 发布在线状态
client.publish('device/status', JSON.stringify({
clientId: CLIENT_ID,
status: 'online',
timestamp: Date.now()
}), { qos: 1, retain: true })
// 模拟发布传感器数据
let count = 0
setInterval(() => {
const payload = JSON.stringify({
clientId: CLIENT_ID,
temperature: 20 + Math.random() * 10,
humidity: 50 + Math.random() * 20,
count: ++count,
timestamp: Date.now()
})
client.publish('sensor/data', payload, { qos: 0 })
console.log(`[${count}] 已发布传感器数据`)
}, 3000)
})
// 收到消息
client.on('message', (topic, payload) => {
const msg = payload.toString()
console.log(`\n[收到消息] 主题: ${topic}`)
console.log(` 内容: ${msg.substring(0, 100)}`)
})
// 断线重连
client.on('reconnect', () => {
console.log('正在重连...')
})
// 错误处理
client.on('error', (err) => {
console.error('MQTT 错误:', err.message)
})
11. MQTT 与其他协议对比
11.1 MQTT vs HTTP
| 对比维度 | MQTT | HTTP |
|---|---|---|
| 模式 | 发布/订阅 | 请求/响应 |
| 协议开销 | 最小 2 字节头部 | 通常 200~800 字节头部 |
| 实时性 | 推送(实时) | 轮询(有延迟) |
| 双向通信 | 原生支持 | 需要 WebSocket |
| 消息确认 | 支持(QoS 1/2) | 不支持 |
| 持久连接 | 长连接 | 短连接(HTTP/1.1) |
| 功耗 | 极低 | 高 |
| 复杂度 | 中低 | 中 |
| 适用场景 | IoT、实时推送 | Web 应用、REST API |
11.2 MQTT vs WebSocket
| 对比维度 | MQTT | WebSocket |
|---|---|---|
| 用途 | 应用层消息协议 | 传输层通信协议 |
| 消息模型 | 发布/订阅 | 点对点 |
| 协议开销 | 小 | 较 MQTT 大 |
| QoS | 原生支持 | 无(依赖上层实现) |
| 主题路由 | 内置 | 无 |
| 保留消息 | 支持 | 不支持 |
| 结合使用 | 可用 WebSocket 传输 MQTT 报文 |
11.3 MQTT vs AMQP
| 对比维度 | MQTT | AMQP |
|---|---|---|
| 定位 | 轻量级 IoT 协议 | 企业级消息中间件协议 |
| 消息模型 | 发布/订阅(主题) | 多种(队列、主题、交换器等) |
| 头部开销 | 2 字节 | 较大 |
| 复杂度 | 低 | 高 |
| 典型实现 | Mosquitto, EMQX | RabbitMQ, ActiveMQ |
| 适用场景 | 传感器、移动设备 | 金融交易、企业集成 |
| 路由能力 | 主题(固定层级) | 交换器 + 绑定(灵活路由) |
11.4 MQTT vs CoAP
| 对比维度 | MQTT | CoAP |
|---|---|---|
| 传输层 | TCP | UDP |
| 模式 | 发布/订阅 | 请求/响应 + 观察 |
| 特性 | 成熟、功能丰富 | 极轻量、专为受限设备设计 |
| 消息格式 | 二进制(自定义) | 二进制(类 HTTP) |
| IP 穿透 | 需要长连接 | 更适合 NAT 环境 |
| 适用场景 | 中高端物联网设备 | 极端受限的设备(如 Z-Wave 替代) |
12. MQTT 最佳实践
12.1 主题设计最佳实践
# ✅ 好的主题命名
sensor/livingroom/temperature
sensor/livingroom/humidity
building/floor3/room301/temperature
factory/line1/machine5/status
device/{device_id}/telemetry
device/{device_id}/command
# ❌ 不好的主题命名
# 1. 不要以 / 开头
/livingroom/temperature # ❌
# 2. 尽量不要使用冗长无意义的单词
this/is/a/very/long/topic/string/that/wastes/bandwidth/with/the/sensor/data/from/device/12345 # ❌
# 3. 避免在主题中使用可变/动态的顶级结构
{i_change}/{constantly}/messing/things/up # ❌
# 4. 同一项目中保持一致的命名风格
SensorLivingroomTemp # ❌ 大小写混用
sensor.livingroom.temperature # ❌ 分隔符不统一(使用 . 而非 /)
主题设计原则:
- 简洁明确:主题名应简短且有自释性
- 层级合理:将不同的维度放在不同的层级
- 统一分隔符:使用
/作为层级分隔符 - 通配符友好:设计主题时考虑订阅者的通配符使用
- 不包含可变数据:设备 ID 可以放在主题中,但时间戳等不应放在主题中
12.2 QoS 选择策略
| 场景 | 推荐 QoS | 原因 |
|---|---|---|
| 传感器周期性上报(如温度、湿度) | QoS 0 | 丢失一次没关系,下次上报会覆盖 |
| 开关控制命令 | QoS 1 | 允许重复,但必须保证执行 |
| 告警消息 | QoS 1 或 2 | 告警不能丢失 |
| 支付交易 | QoS 2 | 必须不重不漏 |
| 位置追踪 | QoS 0 | 高频更新,丢失可接受 |
| 固件升级指令 | QoS 2 | 必须确保恰好执行 |
12.3 安全性最佳实践
# TLS 连接示例
import ssl
# 创建 TLS 上下文
ssl_context = ssl.create_default_context()
ssl_context.check_hostname = True
ssl_context.load_verify_locations("ca.crt") # CA 证书
# 双向 TLS 需要客户端证书
ssl_context.load_cert_chain("client.crt", "client.key")
# MQTT TLS 连接
client.tls_set_context(ssl_context)
client.connect("broker.example.com", 8883, keepalive=60)
安全 Checklist:
- 生产环境使用 TLS 加密(端口 8883)
- 禁止匿名访问
- 使用强密码或证书认证
- 实施 ACL 访问控制
- 限制通配符订阅权限
- 设置合理的 Keep Alive 超时
- 监控异常连接行为
- 定期更新 Broker 和客户端
12.4 性能调优
| 优化项 | 建议 |
|---|---|
| Keep Alive | 根据网络状况设置合理的值(通常 30~120 秒) |
| Payload 大小 | 尽量压缩,大量数据建议使用 Protobuf 或 CBOR 编码 |
| 主题别名 | 长主题场景启用(MQTT 5.0),减少带宽消耗 |
| Receive Maximum | 根据客户端处理能力设置并发消息数 |
| 消息过期 | 设置合理的过期时间,避免无效消息堆积 |
| 批量发布 | 相同主题的高频消息可以合并为一次发布 |
| 连接池 | 服务端使用连接池减少 TCP 握手开销 |
12.5 常见问题与解决方案
| 问题 | 可能原因 | 解决方案 |
|---|---|---|
| 客户端频繁断连 | Keep Alive 太短 / 网络不稳定 | 适当增加 Keep Alive 值,启用自动重连 |
| 消息丢失 | 使用了 QoS 0 + 网络不稳定 | 改用 QoS 1 或 QoS 2 |
| 消息重复 | QoS 1 的重复送达 | 应用层实现去重(使用消息 ID) |
| 订阅者收不到消息 | 主题不匹配 / ACL 限制 | 检查订阅主题的通配符和权限 |
| 连接超时 | Broker 地址错误 / 防火墙拦截 | 检查网络连通性,确认端口放行 |
| Broker 内存暴涨 | 会话未清理 / 保留消息过多 | 设置 Session Expiry,清理无用保留消息 |
| 性能下降 | 订阅数过多 / Payload 过大 | 优化主题层级,压缩 Payload |
13. 附录:MQTT 速查表
13.1 报文类型速查
┌──────────┬─────────┬────────────────┐
│ 报文类型 │ 十六进制 │ 固定报头首字节 │
├──────────┼─────────┼────────────────┤
│ CONNECT │ 1 │ 0x10 │
│ CONNACK │ 2 │ 0x20 │
│ PUBLISH │ 3 │ 0x3x │
│ PUBACK │ 4 │ 0x40 │
│ PUBREC │ 5 │ 0x50 │
│ PUBREL │ 6 │ 0x62 │
│ PUBCOMP │ 7 │ 0x70 │
│ SUBSCRIBE│ 8 │ 0x82 │
│ SUBACK │ 9 │ 0x90 │
│ UNSUBSCRI│ 10 │ 0xA2 │
│ UNSUBACK │ 11 │ 0xB0 │
│ PINGREQ │ 12 │ 0xC0 │
│ PINGRESP │ 13 │ 0xD0 │
│ DISCONN │ 14 │ 0xE0 │
│ AUTH │ 15 │ 0xF0 │
└──────────┴─────────┴────────────────┘
13.2 常用原因码速查
| 原因码 | 名称 | 说明 |
|---|---|---|
| 0x00 | Success | 成功 |
| 0x01 | Normal Disconnection | 正常断开 |
| 0x04 | Disconnect with Will Message | 断开并发送遗嘱 |
| 0x10 | No Matching Subscribers | 无匹配订阅者 |
| 0x11 | No Subscription Existed | 订阅不存在 |
| 0x18 | Continue Authentication | 继续认证 |
| 0x19 | Re-authenticate | 重新认证 |
| 0x80 | Unspecified Error | 未指定错误 |
| 0x81 | Malformed Packet | 报文格式错误 |
| 0x82 | Protocol Error | 协议错误 |
| 0x83 | Implementation Specific Error | 实现特定错误 |
| 0x84 | Unsupported Protocol Version | 不支持的协议版本 |
| 0x85 | Client Identifier Not Valid | 客户端 ID 无效 |
| 0x86 | Bad User Name or Password | 用户名或密码错误 |
| 0x87 | Not Authorized | 未授权 |
| 0x88 | Server Unavailable | 服务端不可用 |
| 0x89 | Server Busy | 服务端繁忙 |
| 0x8A | Banned | 客户端被禁止 |
| 0x8B | Server Shutting Down | 服务端关闭 |
| 0x8C | Bad Authentication Method | 认证方法错误 |
| 0x8D | Keep Alive Timeout | 心跳超时 |
| 0x8E | Session Taken Over | 会话被接管 |
| 0x8F | Topic Filter Invalid | 主题过滤器无效 |
| 0x90 | Topic Name Invalid | 主题名无效 |
| 0x91 | Packet Identifier In Use | 报文标识符已被使用 |
| 0x92 | Packet Identifier Not Found | 报文标识符未找到 |
| 0x93 | Receive Maximum Exceeded | 超出最大接收量 |
| 0x94 | Topic Alias Invalid | 主题别名无效 |
| 0x95 | Packet Too Large | 报文过大 |
| 0x96 | Message Rate Too High | 消息速率过高 |
| 0x97 | Quota Exceeded | 超出配额 |
| 0x98 | Administrative Action | 管理操作 |
| 0x99 | Payload Format Invalid | 载荷格式无效 |
| 0x9A | Retain Not Supported | 不支持保留消息 |
| 0x9B | QoS Not Supported | 不支持的 QoS |
| 0x9C | Use Another Server | 使用其他服务端 |
| 0x9D | Server Moved | 服务端已迁移 |
| 0x9E | Shared Subscriptions Not Supported | 不支持共享订阅 |
| 0x9F | Connection Rate Exceeded | 连接速率超限 |
| 0xA0 | Maximum Connect Time | 最大连接时间 |
| 0xA1 | Subscription IDs Not Supported | 不支持订阅 ID |
| 0xA2 | Wildcard Subscriptions Not Supported | 不支持通配符订阅 |
13.3 常用 MQTT Broker 对比
| Broker | 语言 | 集群 | 插件 | MQTT 5.0 | 管理界面 | 许可证 |
|---|---|---|---|---|---|---|
| EMQX | Erlang | ✅ | ✅ | ✅ | ✅ | 开源 + 商业 |
| Mosquitto | C | ❌ | ❌ | ✅ | 第三方 | EPL-2.0 |
| NanoMQ | C | ✅ | ✅ | ✅ | ✅ | MIT |
| VerneMQ | Erlang | ✅ | ✅ | ✅ | ✅ | Apache 2.0 |
| HiveMQ | Java | ✅ | ✅ | ✅ | ✅ | 商业 |
| RabbitMQ | Erlang | ✅ | ✅ | 插件 | ✅ | MPL 2.0 |
13.4 常用工具
| 工具 | 用途 | 链接 |
|---|---|---|
| MQTTX | 跨平台 MQTT 桌面客户端 | https://mqttx.app |
| MQTT Explorer | MQTT 主题浏览器 | https://mqtt-explorer.com |
| Wireshark | 网络抓包分析 MQTT 报文 | https://wireshark.org |
| mosquitto_pub | 命令行发布工具 | 随 Mosquitto 安装 |
| mosquitto_sub | 命令行订阅工具 | 随 Mosquitto 安装 |
| EMQX Dashboard | EMQX Web 管理界面 | 随 EMQX 安装 |
参考资料:MQTT OASIS 标准规范 v5.0、Eclipse Paho 客户端文档、EMQX 官方文档、Mosquitto 官方手册。
公共测试 Broker:
broker.emqx.io:1883(无加密)/broker.emqx.io:8883(TLS 加密),可用于测试和学习。
编写日期:2026 年 6 月
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐




所有评论(0)