MQTT 协议详解

MQTT(Message Queuing Telemetry Transport)—— 物联网时代最主流的轻量级消息传输协议


目录

  1. MQTT 概述
  2. MQTT 协议核心概念
  3. MQTT 报文格式详解
  4. 各类型报文结构详解
  5. MQTT 5.0 新特性
  6. MQTT 服务质量(QoS)机制
  7. MQTT 会话与持久化
  8. MQTT 安全机制
  9. MQTT 应用场景
  10. MQTT 示例代码
  11. MQTT 与其他协议对比
  12. MQTT 最佳实践
  13. 附录: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/temperaturehome/bedroom/temperature
#(多级通配符) 匹配剩余所有层级 home/# 匹配 home/livingroom/temperaturehome/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   # ❌ 分隔符不统一(使用 . 而非 /)

主题设计原则:

  1. 简洁明确:主题名应简短且有自释性
  2. 层级合理:将不同的维度放在不同的层级
  3. 统一分隔符:使用 / 作为层级分隔符
  4. 通配符友好:设计主题时考虑订阅者的通配符使用
  5. 不包含可变数据:设备 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 官方手册。

公共测试 Brokerbroker.emqx.io:1883(无加密)/ broker.emqx.io:8883(TLS 加密),可用于测试和学习。

编写日期:2026 年 6 月

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐