MQTT 学习笔记

一份覆盖完整知识点的 MQTT 协议学习笔记。重点讲"为什么这么设计"(费曼式 📌 讲解),配代码示例与实战场景。

怎么用这份笔记

  1. 学习/复习 → 顺序看正文,重点看 📌 类比和"为什么"
  2. 查协议 → 翻对应章节,报文结构和参数可直接参考
  3. 实战 → 看 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 使用)

2.2 固定头部(Fixed Header)

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:实时通知、缓存失效广播、低延迟 :::
本页目录