MQTT 学习笔记
一份覆盖完整知识点的 MQTT 协议学习笔记。重点讲"为什么这么设计"(费曼式 📌 讲解),配代码示例与实战场景。
怎么用这份笔记
- 学习/复习 → 顺序看正文,重点看 📌 类比和"为什么"
- 查协议 → 翻对应章节,报文结构和参数可直接参考
- 实战 → 看 Phase 4 的 Broker 部署和客户端编程示例
1 · MQTT 概述与设计理念
1.1 MQTT 是什么
- MQTT(Message Queuing Telemetry Transport)是一种轻量级的发布/订阅消息协议,由 IBM 的 Andy Stanford-Clark 和 Arlen Nipper 于 1999 年发明
- 核心目标:在低带宽、不稳定网络上可靠地传输消息,专为物联网(IoT)和机器间通信(M2M)设计
- 协议版本:MQTT 3.1(2010)、MQTT 3.1.1(2014,最广泛使用)、MQTT 5.0(2019,最新标准)
- 开源协议:MQTT 协议本身免费开放,OASIS 标准化组织维护
TIP
📌 打个比方:MQTT 像一个广播电台系统。Publisher(发布者)是电台节目制作人,Broker(代理)是发射塔,Subscriber(订阅者)是听众。制作人不需要知道谁在听,听众也不需要知道谁在播——发射塔负责把节目推给所有订阅了该频道的听众。这就是发布/订阅解耦。
1.2 为什么选择 MQTT
| 维度 |
MQTT |
HTTP |
AMQP |
CoAP |
| 通信模式 |
发布/订阅 |
请求/响应 |
发布/订阅 |
请求/响应 |
| 协议开销 |
极小(2 字节头部) |
较大(文本头部) |
较大(多帧结构) |
小(4 字节头部) |
| 连接模型 |
长连接 |
短连接为主 |
长连接 |
UDP 无连接 |
| QoS 支持 |
3 级 |
无 |
多级 |
确认重传 |
| 适用场景 |
IoT/传感器 |
Web API |
企业消息 |
受限设备 |
| 传输层 |
TCP |
TCP |
TCP |
UDP |
| 复杂度 |
低 |
低 |
高 |
中 |
TIP
📌 为什么 IoT 首选 MQTT 而不是 HTTP?HTTP 是请求-响应模式——传感器要上报数据得主动发请求,服务器不能主动推送。MQTT 是发布-订阅模式——设备只需建立一次连接,之后双向实时推送。而且 MQTT 头部最小仅 2 字节,HTTP 头部动辄几百字节——在 2G 网络下传 1KB 数据,MQTT 的协议开销可以忽略,HTTP 的头部可能比数据本身还大。
1.3 核心架构
┌──────────┐ PUBLISH ┌─────────┐ SUBSCRIBE ┌──────────────┐
│ Publisher │ ──────────────→ │ Broker │ ──────────────→ │ Subscriber │
│ (传感器) │ │ (代理) │ │ (后端服务) │
└──────────┘ └─────────┘ └──────────────┘
│ │ │
│ CONNECT (TCP 长连接) │ CONNECT (TCP 长连接) │
└───────────────────────────┴──────────────────────────────┘
| 角色 |
作用 |
类比 |
| Publisher |
发布消息到 Topic |
电台节目制作人 |
| Broker |
接收消息并路由给订阅者 |
发射塔 / 中转站 |
| Subscriber |
订阅 Topic 接收消息 |
听众 |
| Topic |
消息的主题/频道 |
电台频道频率 |
| Client |
Publisher 或 Subscriber 的统称 |
对讲机 |
TIP
📌 为什么需要 Broker 中间人?如果设备和后端直接连接,每台设备都要知道后端 IP、处理网络断连重试、维护多后端负载均衡。Broker 把这些复杂性集中处理——设备只需连 Broker,后端也只需连 Broker,Broker 负责路由。这是经典的中介者模式:把 N×M 的网状连接简化为 N+M 的星型连接。
2 · 协议基础
2.1 报文类型
MQTT 定义了 15 种报文类型(MQTT 3.1.1),每种用固定头部第一个字节的 4 位表示:
| 类型 |
ID |
方向 |
说明 |
CONNECT |
1 |
C→B |
客户端请求连接 |
CONNACK |
2 |
B→C |
连接确认 |
PUBLISH |
3 |
C→B / B→C |
发布消息 |
PUBACK |
4 |
C→B / B→C |
QoS 1 发布确认 |
PUBREC |
5 |
C→B / B→C |
QoS 2 发布收到 |
PUBREL |
6 |
C→B / B→C |
QoS 2 发布释放 |
PUBCOMP |
7 |
C→B / B→C |
QoS 2 发布完成 |
SUBSCRIBE |
8 |
C→B |
订阅请求 |
SUBACK |
9 |
B→C |
订阅确认 |
UNSUBSCRIBE |
10 |
C→B |
取消订阅 |
UNSUBACK |
11 |
B→C |
取消订阅确认 |
PINGREQ |
12 |
C→B |
心跳请求 |
PINGRESP |
13 |
B→C |
心跳响应 |
DISCONNECT |
14 |
C→B |
断开连接 |
RESERVED |
15 |
— |
保留(MQTT 5.0 使用) |
Byte 1: ┌─┬─┬─┬─┬─┬─┬─┬─┐
│ 报文类型 (4bit) │ DUP │ QoS │ RETAIN │
└─┴─┴─┴─┴─┴─┴─┴─┴─┘
Byte 2+: 剩余长度(Remaining Length,变长编码)
| 字段 |
位 |
说明 |
| DUP |
3 |
重复标志:是否为重发的消息 |
| QoS |
2-1 |
服务质量等级:0、1、2 |
| RETAIN |
0 |
保留标志:是否为 Retained 消息 |
2.3 剩余长度编码(变长整数)
MQTT 使用变长编码表示剩余长度,最少 1 字节,最多 4 字节:
| 长度范围 |
字节数 |
说明 |
| 0 - 127 |
1 |
最高位 0 |
| 128 - 16383 |
2 |
最高位 1,继续读 |
| 16384 - 2097151 |
3 |
3 字节 |
| 2097152 - 268435455 |
4 |
最大 256MB |
TIP
📌 为什么用变长编码?IoT 设备消息通常很小(几十字节到几 KB),1 字节就能表示 0-127 的长度。如果用固定 4 字节表示长度,每条消息多 3 字节开销——对于每秒发一条温度数据的传感器,一天多 260KB,在 2G 网络下不可忽略。变长编码让小消息用 1 字节,大消息用最多 4 字节,按需付费。
3 · QoS 服务质量
3.1 三级 QoS
| QoS |
名称 |
机制 |
保证 |
开销 |
| 0 |
至多一次 (At most once) |
发完即忘,不确认 |
消息可能丢失 |
最小 |
| 1 |
至少一次 (At least once) |
PUBACK 确认,可能重复 |
消息不会丢,可能重复 |
中等 |
| 2 |
恰好一次 (Exactly once) |
四步握手 (PUBREC→PUBREL→PUBCOMP) |
不丢不重 |
最大 |
3.2 QoS 0 — 至多一次
Sender Receiver
│ ── PUBLISH ──→ │
│ │
│ (不等待确认) │
- 消息发出后不等待任何确认,发完即忘
- 网络丢包则消息丢失,不重发
- 适用于高频且允许丢失的数据(如温度传感器每秒上报,丢几条无所谓)
3.3 QoS 1 — 至少一次
Sender Receiver
│ ── PUBLISH ──→ │
│ │
│ ←── PUBACK ── │
│ │
│ (超时未收到 PUBACK │
│ 则重发,DUP=1) │
- 发送方存储消息,等待 PUBACK
- 超时未收到则重发,DUP 标志置 1
- 接收方可能收到重复消息——业务层需做幂等处理
- 适用于不能丢但允许重复的消息(如设备状态变更通知)
TIP
📌 QoS 1 为什么会重复?发送方发 PUBLISH,接收方收到并回 PUBACK。但 PUBACK 在网络中丢了——发送方超时后重发 PUBLISH(DUP=1),接收方再次收到同一条消息。协议层面无法消除这种重复,因为接收方无法区分"新消息"和"PUBACK 丢失后的重发"。去重是应用层的责任。
3.4 QoS 2 — 恰好一次
Sender Receiver
│ ── PUBLISH ──→ │
│ │
│ ←── PUBREC ── │ (收到,记录 PacketID)
│ │
│ ── PUBREL ──→ │ (可以投递了)
│ │
│ ←── PUBCOMP ── │ (完成)
- 四步握手确保消息不丢不重
- 接收方收到 PUBLISH 后存 PacketID,回 PUBREC
- 发送方收到 PUBREC 后发 PUBREL,接收方收到后才投递消息并回 PUBCOMP
- 适用于不能丢也不能重复的消息(如计费指令、支付命令)
TIP
📌 QoS 2 为什么需要四步而不是两步?两步(PUBLISH + PUBACK)就是 QoS 1,会重复。QoS 2 的思路是:先确认"我收到了"(PUBREC),再确认"你可以处理了"(PUBREL),最后确认"我处理完了"(PUBCOMP)。接收方在收到 PUBREL 之前不投递消息——如果 PUBLISH 重复到达,接收方发现 PacketID 已存在,只回 PUBREC 不重复投递。用两轮确认换取恰好一次。
3.5 QoS 选择建议
| 场景 |
推荐 QoS |
理由 |
| 温度/湿度传感器 |
0 |
高频、容错、省带宽 |
| 设备上下线通知 |
1 |
不能丢,重复无害 |
| GPS 位置上报 |
1 |
不能丢,重复可去重 |
| 支付/计费指令 |
2 |
不能丢不能重复 |
| 日志采集 |
0 或 1 |
看日志重要性 |
| 远程控制开关 |
1 |
不能丢,重复无害(幂等操作) |
TIP
📌 QoS 越高越好吗?不是。QoS 2 的四步握手意味着每条消息 4 次网络往返——在 200ms 延迟的 2G 网络下,一条消息要 800ms。QoS 0 只需 1 次发送。按业务需求选最低够用的 QoS,不要无脑用 QoS 2。
4 · Topic 与通配符
4.1 Topic 层级
Topic 用 / 分隔层级,类似文件路径:
home/livingroom/temperature
home/livingroom/humidity
home/bedroom/temperature
factory/line1/sensor3/vibration
TIP
📌 为什么 Topic 用层级而不是扁平名?层级结构天然支持过滤和分组。订阅 home/+/temperature 一次获取所有房间的温度,订阅 home/# 获取家里所有数据。扁平命名(如 home_livingroom_temperature)做不到这种模式匹配。层级 = 结构化 = 可聚合。
4.2 通配符
| 通配符 |
匹配规则 |
示例 |
匹配 |
+ |
匹配单层 |
home/+/light |
home/kitchen/light ✓home/bedroom/light ✓home/kitchen/fan/light ✗ |
# |
匹配多层(必须放最后) |
home/# |
home/kitchen/temp ✓home/bedroom/ac/fan ✓office/temp ✗ |
4.3 Topic 设计原则
# 好的设计:层级清晰、可聚合
sensor/floor3/room301/temperature
sensor/floor3/room301/humidity
device/abc123/status
device/abc123/command
# 坏的设计:扁平、无法聚合
floor3_room301_temperature
device_abc123_status_get
| 原则 |
说明 |
| 层级从粗到细 |
building/floor/room/sensor |
不要以 / 开头 |
/home/temp 会被解析为空层级 + home + temp |
| 避免空层级 |
home//temp 中间空层级容易出错 |
| 使用小写 |
Topic 大小写敏感,统一小写避免混乱 |
| 不要用空格和特殊字符 |
只用字母、数字、/、_、- |
| 设备 ID 放中间层 |
device/{id}/sensor/{type} 而非 {id}/sensor/{type} |
TIP
📌 Topic 设为什么会影响性能吗?会。Broker 对 Topic 做树形索引匹配——层级越深,匹配开销越大。但这个开销在万级 Topic 以下可以忽略。真正影响性能的是订阅数量和消息频率,不是 Topic 层级深度。设计 Topic 优先考虑可读性和可聚合性,不要为了性能牺牲结构。
5 · 连接与心跳
5.1 CONNECT 报文
CONNECT 报文包含:
├── Client ID (客户端唯一标识)
├── Username (可选,认证用户名)
├── Password (可选,认证密码)
├── Clean Session (是否清除会话)
├── Keep Alive (心跳间隔,秒)
├── Will Topic (遗嘱主题,可选)
├── Will Message (遗嘱消息,可选)
├── Will QoS (遗嘱 QoS,可选)
├── Will Retain (遗嘱是否保留)
└── Payload (应用消息,可选)
5.2 Keep Alive 心跳机制
Client Broker
│ ── CONNECT (KeepAlive=60) ──→ │
│ ←── CONNACK ── │
│ │
│ (60秒内无消息交换) │
│ ── PINGREQ ──→ │
│ ←── PINGRESP ── │
│ │
│ (如果 1.5×60=90秒无响应) │
│ 判定连接断开 │
- 客户端在
Keep Alive 秒内无消息交换时发送 PINGREQ
- Broker 回
PINGRESP
- 如果 Broker 在 1.5 × Keep Alive 时间内未收到任何消息,判定客户端离线
- Keep Alive = 0 表示不使用心跳(不推荐)
TIP
📌 为什么需要心跳?TCP 连接断开有两种:正常关闭(FIN)和异常断开(断电、信号丢失)。异常断开时 TCP 不会发送 FIN——对端不知道连接已断,直到尝试写入时收到 RST。心跳机制让双方在空闲时也能检测连接状态,及时发现异常断开并触发遗嘱消息。心跳 = 主动探测 vs 被动等待超时。
5.3 Clean Session
| Clean Session |
行为 |
true |
连接断开后清除所有会话状态(订阅、未送达消息)。重连后是全新会话 |
false |
保留会话状态。重连后恢复订阅,收到离线期间积压的消息 |
TIP
📌 什么时候用 Clean Session = false?设备可能频繁断网(如移动设备进隧道),但离线期间的消息不能丢(如控制指令)。设 false 后,Broker 为该客户端存储离线消息,重连后投递。代价是 Broker 要持久化存储——如果 100 万设备都设 false,每个积压 100 条消息,Broker 要存 1 亿条。按需使用,设合理的过期时间。
6 · 发布与订阅
6.1 发布消息(PUBLISH)
PUBLISH 报文:
├── Topic Name (目标主题)
├── QoS (0/1/2)
├── RETAIN (是否保留)
├── DUP (是否重发)
├── Packet ID (QoS>0 时有)
└── Payload (消息内容,二进制)
- Payload 是原始二进制,可以是 JSON、Protobuf、纯文本、图片——协议不限制
- 最大消息大小由 Broker 配置决定(Mosquitto 默认 0 = 无限制)
6.2 订阅(SUBSCRIBE)
SUBSCRIBE 报文:
├── Packet ID
└── Topic Filters (可同时订阅多个)
├── Topic Filter 1 + QoS
├── Topic Filter 2 + QoS
└── ...
# Python 示例
import paho.mqtt.client as mqtt
client = mqtt.Client()
# 订阅(可指定 QoS)
client.subscribe("home/+/temperature", qos=1)
client.subscribe([("device/+/status", 0), ("alert/#", 2)]) # 批量订阅
# 取消订阅
client.unsubscribe("home/+/temperature")
6.3 订阅的 QoS 降级
消息投递的 QoS = min(发布 QoS, 订阅 QoS):
| 发布 QoS |
订阅 QoS |
实际投递 QoS |
| 2 |
2 |
2 |
| 2 |
1 |
1 |
| 2 |
0 |
0 |
| 1 |
2 |
1 |
| 0 |
2 |
0 |
TIP
📌 为什么取最小值?QoS 是发送方的承诺——QoS 2 的发布者保证"恰好一次",但 QoS 0 的订阅者没做接收确认机制,无法保证。取 min 确保不"虚假承诺":订阅者没声明能处理 QoS 2,就不能强制给它 QoS 2 的保证。QoS 是双方协商的最低公约数。
7 · Retained 消息
7.1 什么是 Retained 消息
- 发布消息时
RETAIN=1,Broker 会保留这条消息
- 新订阅者订阅该 Topic 时,立即收到这条保留消息
- 每个 Topic 只保留最后一条 Retained 消息
# 发布保留消息
client.publish("device/sensor1/config", '{"interval": 5}', retain=True, qos=1)
# 之后任何新订阅者订阅 device/sensor1/config
# 都会立即收到 {"interval": 5}
7.2 删除 Retained 消息
发布一个 Payload 为空 的 Retained 消息即可删除:
client.publish("device/sensor1/config", "", retain=True) # 空payload = 删除
TIP
📌 Retained 消息解决什么问题?新设备上线后需要知道当前配置(如采样间隔、阈值),但配置消息可能早就发过了。没有 Retained,新设备订阅后只能等下一次配置变更才能收到——可能等几小时。有 Retained,Broker 立即把最后一次配置推给新设备。Retained = Topic 的"最后状态快照"。
📌 什么时候不该用 Retained?高频变化的数据(如温度每秒更新)——每次都设 Retained 会让 Broker 不断覆写存储,浪费资源。Retained 适合低频、状态型数据(配置、开关状态、设备型号)。
8 · Last Will 与遗嘱
8.1 什么是遗嘱消息
- 客户端连接时声明一个"遗嘱"——如果异常断开(非正常 DISCONNECT),Broker 自动发布这条遗嘱消息
- 用于通知其他客户端"我掉线了"
client = mqtt.Client()
client.will_set(
topic="device/sensor1/status",
payload="offline",
qos=1,
retain=True
)
client.connect("broker.example.com", 1883, 60)
# 正常断开:不发遗嘱
client.disconnect()
# 异常断开(断电、网络丢失):Broker 在 1.5×KeepAlive 后发遗嘱
# → "device/sensor1/status" = "offline" (retained)
8.2 遗嘱触发条件
| 断开方式 |
是否触发遗嘱 |
| 客户端发 DISCONNECT |
不触发 |
| Keep Alive 超时 |
触发 |
| TCP 连接断开(Broker 检测到) |
触发 |
| Broker 主动断开(协议违规) |
触发 |
| 客户端崩溃 |
触发 |
TIP
📌 为什么正常 DISCONNECT 不触发遗嘱?正常断开是客户端主动告知"我要走了",不需要遗嘱——客户端自己可以发一条"offline"消息。遗嘱是为不可控的异常断开设计的——设备断电、信号丢失时无法主动发消息,Broker 代发。如果正常断开也触发遗嘱,客户端就没法区分"主动离线"和"异常掉线"了。
📌 Retained + Will 的经典组合:设备上线时发 status=online(retained),遗嘱设 status=offline(retained)。任何订阅者随时能读到设备当前在线状态——在线时是 online,异常掉线后 Broker 自动改为 offline。这是 IoT 设备状态管理的标准模式。
9 · MQTT 5.0 新特性
9.1 主要新增
| 特性 |
说明 |
| Reason Code |
CONNACK / DISCONNECT 等携带原因码,不再"哑巴断开" |
| User Properties |
自定义键值对头(类似 HTTP Header) |
| Message Expiry |
消息过期时间,超时自动丢弃 |
| Topic Alias |
用整数别名替代长 Topic 名,省带宽 |
| Shared Subscription |
负载均衡订阅($share/group/topic) |
| Will Delay |
遗嘱延迟发布,避免短暂断线误报 |
| Session Expiry Interval |
会话过期时间,替代 Clean Session 布尔值 |
| Maximum Packet Size |
协商最大包大小 |
| Server Reference |
Broker 重定向到另一个 Broker |
| Subscription Identifier |
订阅 ID,区分消息来自哪个订阅 |
9.2 Reason Code 示例
MQTT 3.1.1: CONNACK 只有 0(成功) 和 1(失败)
MQTT 5.0: CONNACK 有 30+ 原因码
0x00 Success
0x04 Authentication can be retried
0x05 Not authorized
0x14 Unsupported protocol version
0x17 Server unavailable
0x86 Specified topic name too long
0x87 Quota exceeded
TIP
📌 MQTT 5.0 的 Reason Code 为什么重要?MQTT 3.1.1 连接失败只告诉你"失败"——是密码错?协议版本不对?还是 Broker 满了?客户端只能猜。MQTT 5.0 给出具体原因码,客户端可以智能重试——密码错就不重试,Broker 满了可以等一会儿再试。从"哑巴"到"会说话"。
9.3 Topic Alias
# 第一次:发布时声明别名
PUBLISH:
Topic: "factory/line1/machine3/sensor/temperature" (38字节)
Topic Alias: 1
# 之后:只发别名编号
PUBLISH:
Topic: "" (0字节)
Topic Alias: 1
Payload: 25.5
TIP
📌 Topic Alias 省多少?如果 Topic 是 factory/line1/machine3/sensor/temperature(38 字节),每条消息都要带。如果 Payload 只有 25.5(4 字节),Topic 占了 90% 的包大小。用 Alias 后后续消息 Topic 字段为 0 字节,只传一个 2 字节的 Alias ID。高频小消息场景下省带宽效果显著。
10 · Shared Subscription
10.1 普通订阅 vs 共享订阅
# 普通订阅:每个订阅者都收到全量消息
Topic: task/queue
Subscriber A → 收到 task1, task2, task3
Subscriber B → 收到 task1, task2, task3 (重复!)
# 共享订阅:消息在订阅者间负载均衡
Topic: $share/workers/task/queue
Subscriber A → 收到 task1, task3
Subscriber B → 收到 task2 (不重复!)
10.2 语法
$share/{group}/{topic_filter}
# 示例
$share/processor/device/+/event # 同组负载均衡
$share/gpu_processor/render/queue # GPU 渲染队列
{group} 是组名,同组内负载均衡
- 不同组各自收到全量消息(类似广播到各组,组内再分发)
TIP
📌 Shared Subscription 解决什么问题?如果 1000 个设备发消息到 device/event,后端 1 个订阅者处理不过来。加更多订阅者——但普通订阅下每个订阅者都收到全量消息,1000 条消息变成 1000×N 条。Shared Subscription 让 N 个订阅者分摊消息——每人处理 1000/N 条。这是 MQTT 版的消费者组模式,类似 Kafka 的 Consumer Group。
11 · Session 与持久化
11.1 Session 状态
| 存储内容 |
Clean Session=true |
Clean Session=false |
| 订阅列表 |
连接时清除 |
保留 |
| 离线消息(QoS≥1) |
不存储 |
存储并重连后投递 |
| 未确认消息 |
清除 |
保留重传 |
| Packet ID |
重置 |
继续 |
11.2 Broker 持久化
# Mosquitto 持久化配置
persistence true
persistence_location /var/lib/mosquitto/
autosave_interval 1800 # 30分钟自动保存
- 内存型 Broker(默认):重启后所有 Session 丢失
- 持久化型 Broker:Session 状态写入磁盘,重启后恢复
TIP
📌 持久化 vs 内存模式怎么选?如果设备都是 Clean Session=true(不需要离线消息),内存模式即可——Broker 重启后设备重连,全新会话。如果有设备用 Clean Session=false(需要离线消息),必须持久化——否则 Broker 重启后离线消息全丢。持久化有写入开销,按业务需求决定。
12 · 安全机制
12.1 认证方式
| 方式 |
说明 |
安全等级 |
| 匿名 |
无认证(allow_anonymous true) |
仅开发用 |
| 用户名/密码 |
CONNECT 报文携带 |
明文(需 TLS) |
| 客户端证书 |
TLS 双向认证 |
高 |
| Token/OAuth |
MQTT 5.0 增强认证 |
高 |
| 自定义 |
Broker 插件(如 Mosquitto auth-plugin) |
灵活 |
12.2 TLS 加密
# Mosquitto TLS 配置
listener 8883
cafile /etc/mosquitto/certs/ca.crt
certfile /etc/mosquitto/certs/server.crt
keyfile /etc/mosquitto/certs/server.key
require_certificate true # 双向认证
# Python 客户端 TLS
import ssl
client.tls_set(
ca_certs="ca.crt",
certfile="client.crt",
keyfile="client.key",
tls_version=ssl.PROTOCOL_TLS
)
client.connect("broker.example.com", 8883)
12.3 ACL 访问控制
# Mosquitto ACL 配置
user sensor1
topic read device/sensor1/#
topic write device/sensor1/data
user backend
topic read device/#
topic write device/+/command
pattern read $SYS/#
TIP
📌 为什么 MQTT 安全特别重要?MQTT Broker 常暴露在公网(设备要远程连接)。如果无认证——任何人都能订阅你的设备数据或发控制指令。最低限度也要用户名/密码 + TLS。生产环境推荐客户端证书——每台设备一个证书,即使密码泄露没有证书也连不上。
📌 为什么不用 HTTP 的安全方案?MQTT 是长连接,不能用 HTTP 的 Bearer Token 模式(每次请求带 Token)。MQTT 只在 CONNECT 时认证一次,之后整个连接生命周期不再认证。所以认证信息要"够重"——证书比密码更难窃取。
13 · Broker 选型与部署
13.1 主流 Broker 对比
| Broker |
语言 |
特点 |
适用场景 |
| Mosquitto |
C |
轻量、稳定、资源占用极低 |
边缘网关、小型部署 |
| EMQX |
Erlang |
高性能(千万连接)、功能丰富 |
大规模 IoT 平台 |
| HiveMQ |
Java |
企业级、MQTT 5.0 完整支持 |
企业 IoT |
| RabbitMQ |
Erlang |
AMQP 为主,MQTT 插件 |
混合消息队列 |
| VerneMQ |
Erlang |
高性能、插件化 |
中大规模部署 |
13.2 Mosquitto 部署(Docker)
# docker-compose.yml
version: '3'
services:
mosquitto:
image: eclipse-mosquitto:2
ports:
- "1883:1883" # MQTT
- "9001:9001" # WebSocket
volumes:
- ./mosquitto.conf:/mosquitto/config/mosquitto.conf
- mosquitto_data:/mosquitto/data
- mosquitto_log:/mosquitto/log
volumes:
mosquitto_data:
mosquitto_log:
# mosquitto.conf
persistence true
persistence_location /mosquitto/data/
log_dest file /mosquitto/log/mosquitto.log
# 匿名关闭,密码认证
allow_anonymous false
password_file /mosquitto/config/passwd
# 监听
listener 1883
protocol mqtt
# WebSocket(前端直连)
listener 9001
protocol websockets
# ACL
acl_file /mosquitto/config/aclfile
# 生成密码文件
docker exec mosquitto mosquitto_passwd -c /mosquitto/config/passwd user1
# 输入密码...
# 重启生效
docker restart mosquitto
13.3 EMQX 部署(Docker)
version: '3'
services:
emqx:
image: emqx/emqx:5.0
ports:
- "1883:1883" # MQTT
- "8083:8083" # WebSocket
- "8084:8084" # WSS
- "8883:8883" # MQTTS
- "18083:18083" # Dashboard
environment:
EMQX_NAME: emqx
EMQX_HOST: 127.0.0.1
TIP
📌 Mosquitto vs EMQX 怎么选?Mosquitto 是 C 写的单线程——轻量但上限在几万连接。EMQX 是 Erlang 写的——天生支持高并发,单节点百万连接,集群千万级。边缘网关用 Mosquitto,云端平台用 EMQX。如果你不确定,先用 Mosquitto——简单稳定,等规模上来再迁 EMQX。
14 · 客户端编程
14.1 Python(paho-mqtt)
import paho.mqtt.client as mqtt
import json
# 回调函数
def on_connect(client, userdata, flags, rc):
print(f"连接结果: {rc}")
client.subscribe("home/+/temperature", qos=1)
def on_message(client, userdata, msg):
payload = json.loads(msg.payload.decode())
print(f"Topic: {msg.topic}, QoS: {msg.qos}, Data: {payload}")
def on_disconnect(client, userdata, rc):
print(f"断开连接: {rc}")
# 创建客户端
client = mqtt.Client(client_id="my_client")
client.username_pw_set("user1", "password")
client.on_connect = on_connect
client.on_message = on_message
client.on_disconnect = on_disconnect
# 遗嘱
client.will_set("device/my_client/status", "offline", qos=1, retain=True)
# 连接
client.connect("broker.example.com", 1883, 60)
# 发布
client.publish("home/livingroom/temperature",
json.dumps({"value": 23.5, "unit": "C"}),
qos=1, retain=True)
# 上线状态
client.publish("device/my_client/status", "online", retain=True)
# 后台循环
client.loop_start() # 非阻塞
# ... 做其他事 ...
client.loop_stop()
14.2 JavaScript(mqtt.js)
import mqtt from 'mqtt';
const client = mqtt.connect('mqtt://broker.example.com:1883', {
clientId: 'js_client',
username: 'user1',
password: 'password',
keepalive: 60,
clean: false, // 保留会话
will: {
topic: 'device/js_client/status',
payload: 'offline',
qos: 1,
retain: true
}
});
client.on('connect', () => {
console.log('已连接');
client.subscribe('home/+/temperature', { qos: 1 });
client.publish('device/js_client/status', 'online', { retain: true });
});
client.on('message', (topic, payload) => {
const data = JSON.parse(payload.toString());
console.log(`${topic}:`, data);
});
client.on('error', (err) => {
console.error('MQTT 错误:', err);
});
client.on('offline', () => {
console.log('连接断开,自动重连中...');
});
// 发布
client.publish('home/livingroom/light', 'on', { qos: 1 });
14.3 C(嵌入式常用库)
#include "MQTTClient.h"
#define ADDRESS "tcp://broker.example.com:1883"
#define CLIENTID "embedded_device_001"
#define TOPIC "device/001/sensor"
#define QOS 1
#define TIMEOUT 10000L
int main() {
MQTTClient client;
MQTTClient_connectOptions conn_opts = MQTTClient_connectOptions_initializer;
MQTTClient_message pubmsg = MQTTClient_message_initializer;
MQTTClient_deliveryToken token;
MQTTClient_create(&client, ADDRESS, CLIENTID,
MQTTCLIENT_PERSISTENCE_NONE, NULL);
conn_opts.keepAliveInterval = 60;
conn_opts.cleansession = 1;
conn_opts.username = "user1";
conn_opts.password = "password";
MQTTClient_connect(client, &conn_opts);
// 发布
char payload[32];
snprintf(payload, sizeof(payload), "{\"temp\": %.1f}", 25.5);
pubmsg.payload = payload;
pubmsg.payloadlen = strlen(payload);
pubmsg.qos = QOS;
pubmsg.retained = 1;
MQTTClient_publishMessage(client, TOPIC, &pubmsg, &token);
MQTTClient_waitForCompletion(client, token, TIMEOUT);
MQTTClient_disconnect(client, 10000);
MQTTClient_destroy(&client);
return 0;
}
TIP
📌 嵌入式设备选什么库?内存极小(< 64KB RAM)用 Mbed MQTT 或 esp-mqtt(ESP32)。有操作系统(Linux)用 Paho MQTT C。资源够用(> 1MB RAM)可以用 Paho MQTT C++ 或 mosquitto client。选库的核心指标是RAM 占用和代码体积,不是功能多少。
15 · 物联网架构设计
15.1 典型 IoT 架构
┌─────────┐ MQTT ┌──────────┐ MQTT ┌──────────┐
│ 设备层 │ ──────────→ │ 边缘网关 │ ──────────→ │ 云端Broker│
│ (传感器) │ (局域网) │(Mosquitto)│ (公网/TLS) │ (EMQX) │
└─────────┘ └──────────┘ └──────────┘
│
┌─────────────────┼──────────────┐
│ │ │
┌─────┴────┐ ┌──────┴────┐ ┌──────┴────┐
│ 数据存储 │ │ 规则引擎 │ │ 后端服务 │
│(TSDB/DB) │ │(EMQX规则) │ │(API/告警) │
└──────────┘ └──────────┘ └──────────┘
15.2 Topic 设计实例
# 智能家居系统
home/{room}/{device_type}/{action}
home/livingroom/thermostat/temperature # 温度上报
home/livingroom/thermostat/setpoint # 设定温度
home/livingroom/light/status # 灯状态
home/livingroom/light/command # 灯控制指令
home/+/sensor/temperature # 所有房间温度
home/# # 全屋所有消息
# 工厂设备监控
factory/{line}/{machine}/{sensor}/{metric}
factory/line1/machine3/vibration/rms
factory/line1/machine3/temperature/value
factory/line1/machine3/status # 运行/停机/故障
factory/+/machine3/# # 产线1的3号机全数据
factory/line1/+/status # 产线1所有设备状态
15.3 边缘网关模式
传感器1 ─┐
传感器2 ─┼─→ 边缘 Mosquitto ──TLS──→ 云端 EMQX
传感器3 ─┘ (局域网, 无TLS) (公网, 双向TLS)
- 边缘网关聚合局域网内设备数据
- 设备连本地网关(低延迟、无公网暴露)
- 网关桥接(Bridge)到云端 Broker(TLS 加密)
# Mosquitto Bridge 配置
connection bridge-to-cloud
address cloud-broker.example.com:8883
topic home/# both 1
bridge_cafile /etc/mosquitto/certs/ca.crt
bridge_certfile /etc/mosquitto/certs/client.crt
bridge_keyfile /etc/mosquitto/certs/client.key
bridge_attempt_unsubscribe false
TIP
📌 为什么需要边缘网关?工厂里 200 个传感器如果各自连云端 Broker——200 条公网 TLS 连接,每条握手 10KB,200 个证书管理。用边缘网关:200 个传感器连本地 Mosquitto(无 TLS,毫秒延迟),Mosquitto 用 1 条 TLS 连接桥接到云端。聚合连接、降低延迟、简化证书管理。
16 · 监控与运维
16.1 $SYS 主题
Broker 内置 $SYS/ 主题暴露运行指标:
$SYS/broker/uptime # 运行时间
$SYS/broker/clients/connected # 当前连接数
$SYS/broker/clients/total # 总连接数
$SYS/broker/messages/received # 接收消息总数
$SYS/broker/messages/sent # 发送消息总数
$SYS/broker/messages/stored # 存储消息数
$SYS/broker/subscriptions/count # 订阅数
$SYS/broker/retained messages/count # 保留消息数
# 监控连接数
mosquitto_sub -t '$SYS/broker/clients/connected' -h localhost
16.2 常见运维问题
| 问题 |
排查方法 |
| 设备连不上 |
检查认证、TLS 证书、防火墙端口 |
| 消息延迟高 |
检查 QoS 等级、消息大小、订阅数 |
| Broker 内存高 |
检查离线消息积压、Retained 消息数 |
| 消息丢失 |
检查 QoS 等级、Clean Session 设置 |
| 消息重复 |
QoS 1 正常现象,业务层做幂等 |
| 遗嘱没触发 |
检查 Keep Alive 设置、是否正常 DISCONNECT |
| 连接频繁断开 |
检查 Keep Alive 间隔是否太短 |
16.3 命令行工具
# mosquitto_pub 发布
mosquitto_pub -h broker.example.com -p 1883 \
-t "test/topic" -m "hello" -q 1 -r \
-u user1 -P password
# mosquitto_sub 订阅
mosquitto_sub -h broker.example.com -p 1883 \
-t "test/#" -q 1 \
-u user1 -P password -v
# TLS 连接
mosquitto_pub -h broker.example.com -p 8883 \
--cafile ca.crt --cert client.crt --key client.key \
-t "test/topic" -m "hello"
附录 A · 命令速查
mosquitto_pub
| 参数 |
作用 |
-h |
Broker 主机 |
-p |
端口 |
-t |
Topic |
-m |
消息内容 |
-q |
QoS (0/1/2) |
-r |
Retained 消息 |
-u / -P |
用户名/密码 |
-f |
从文件读取 Payload |
--cafile |
CA 证书 |
--cert / --key |
客户端证书/私钥 |
mosquitto_sub
| 参数 |
作用 |
-t |
订阅 Topic |
-T |
排除 Topic |
-v |
显示 Topic 前缀 |
-C |
收到 N 条后退出 |
-W |
超时秒数 |
-c |
Clean Session = false |
附录 B · 报文结构速查
固定头部
┌───────────────┬───────┬───────┬────────┐
│ 报文类型 (4bit) │ DUP │ QoS │ RETAIN │
└───────────────┴───────┴───────┴────────┘
┌──────────────────────────────────────────┐
│ 剩余长度 (1-4 字节变长编码) │
└──────────────────────────────────────────┘
CONNECT 报文
固定头部: 类型=1 + 剩余长度
可变头部:
Protocol Name (MQTT / MQTT5)
Protocol Level (4=3.1.1, 5=5.0)
Connect Flags (1字节: UserFlag, PasswordFlag, WillRetain,
WillQoS(2bit), WillFlag, CleanSession, Reserved)
Keep Alive (2字节)
Payload:
Client ID
Will Topic + Will Message (if WillFlag=1)
Username + Password (if flags=1)
PUBLISH 报文
固定头部: 类型=3 + DUP + QoS + RETAIN + 剩余长度
可变头部:
Topic Name (UTF-8 字符串)
Packet ID (仅 QoS>0)
Payload:
应用消息 (二进制, 0 到 N 字节)
附录 C · 常见问题 FAQ
Q: MQTT 和 WebSocket 什么关系?
A: 两个不同协议。MQTT 是应用层协议跑在 TCP 上。但浏览器不能直接用 TCP——所以 Broker 提供 WebSocket 端口(如 9001),前端用 mqtt.js 通过 WebSocket 传输 MQTT 报文。MQTT over WebSocket = 让浏览器也能用 MQTT。
Q: 一个 Broker 能支持多少连接?
A: Mosquitto 单节点约 10 万级。EMQX 单节点百万级,集群千万级。瓶颈通常是内存(每个 TCP 连接约 4-50KB)和文件描述符限制(ulimit -n)。
Q: MQTT 消息能保证顺序吗?
A: 同一 Topic 同一 QoS 的消息按发布顺序投递。但不同 Topic 之间无顺序保证。跨 Broker 桥接也不保证顺序。
Q: 消息最大能多大?
A: 协议理论上限 256MB(4 字节变长编码最大值)。实际由 Broker 配置限制(Mosquitto 默认 0=无限,EMQX 默认 1MB)。IoT 场景建议不超过 256KB。
Q: 如何实现请求-响应模式?
A: MQTT 本身是发布/订阅,不直接支持请求-响应。常见做法:客户端发布请求到 device/{id}/command,并在 device/{id}/response 上订阅响应。MQTT 5.0 新增了 Response Topic 和 Correlation Data 简化此模式。
Q: Broker 重启后 Retained 消息会丢吗?
A: 如果 Broker 开启了持久化(persistence true),Retained 消息会保存到磁盘,重启后恢复。如果内存模式,重启后丢失。
附录 D · 与其他消息协议对比
| 维度 |
MQTT |
Kafka |
RabbitMQ |
AMQP |
Redis Pub/Sub |
| 定位 |
IoT 消息协议 |
流处理平台 |
消息队列 |
企业消息标准 |
内存发布/订阅 |
| 协议 |
MQTT |
自定义 |
AMQP/MQTT |
AMQP |
RESP |
| QoS |
3 级 |
分区有序 |
多种确认 |
多种确认 |
无 |
| 持久化 |
可选 |
强持久化 |
可选 |
可选 |
无 |
| 消息大小 |
小(KB级) |
大(MB级) |
中 |
中 |
小 |
| 连接数 |
百万级 |
万级 |
万级 |
万级 |
万级 |
| 延迟 |
毫秒 |
毫秒 |
毫秒 |
毫秒 |
微秒 |
| 适用 |
IoT/传感器 |
日志/流处理 |
企业应用 |
企业集成 |
缓存/实时 |
:::tip
📌 选型建议:
- MQTT:IoT 设备通信、低带宽网络、海量连接
- Kafka:高吞吐日志流、事件溯源、数据管道
- RabbitMQ:企业应用集成、复杂路由、任务队列
- Redis Pub/Sub:实时通知、缓存失效广播、低延迟
:::