MQTT 物联网消息协议 — 小白讲解 + 代码实战 + 面试题
MQTT 物联网消息协议 — 小白讲解 + 代码实战 + 面试题
定位:假设你只用过 HTTP 和 RabbitMQ,没接触过物联网。所有概念从零讲起,配可运行代码 + 生产踩坑 + 面试题。 覆盖:MQTT 是什么 → 核心机制(QoS / 主题 / 保留消息 / 遗嘱 / 会话)→ 报文结构 → 代码实战(Java / Spring / Python / 前端)→ 应用场景 → 15 个生产问题与解决方案 → 安全 → 选型对比 → 面试题。 用法:先读讲解理解概念 → 用 Docker 起个 Broker 跑代码验证 → 再看面试题自测。 特色:★ 每个难概念都按「一句话定义 → 生活类比 → 时序图 → 致命缺点 → 什么时候用」五步讲;每个场景配可直接运行的代码。
📖 本文档名词速查(看到不认识的缩写,先查这里)
完整的白话解释、生活类比、面试话术见
00-名词速查手册.md。 下面只列本文档用到的名词,按在文档中出现的顺序排序。
| 缩写 / 术语 | 英文全称 | 一句话说明 | 出处 |
|---|---|---|---|
| MQTT | Message Queuing Telemetry Transport | 物联网用的超轻量消息协议,发布/订阅模式,为弱网和小设备设计 | 第2章 |
| Broker | Broker(代理服务器) | MQTT 的服务端,所有消息的中转站(相当于“邮局”) | 第2章 |
| 发布/订阅 | Publish / Subscribe(Pub/Sub) | 发消息的人不知道谁会收到;收消息的人订阅感兴趣的频道 | 第2章 |
| Topic | Topic(主题) | 消息的“频道名”,用 / 分层,如 home/livingroom/temp |
第2章 |
| QoS | Quality of Service | 消息投递的可靠等级:0 最多一次 / 1 至少一次 / 2 恰好一次 | 第2章 |
| 发布者 | Publisher | 发消息的一方(设备或应用) | 第2章 |
| 订阅者 | Subscriber | 收消息的一方 | 第2章 |
| 保留消息 | Retained Message | Broker 记住每个主题的最后一条消息,新订阅者上线立刻收到 | 第2章 |
| 遗嘱消息 | Last Will and Testament(LWT) | 设备提前留好的“遗言”,它意外掉线时 Broker 自动代发 | 第2章 |
| 会话 | Session | Broker 为客户端保留的“档案”(订阅关系、未确认的消息) | 第2章 |
| Clean Session | Clean Session(MQTT 3.1.1) | 连接时设 true = 不保留档案,断开就忘;false = 保留 |
第2章 |
| Session Expiry | Session Expiry Interval(MQTT 5.0) | 5.0 里的会话保留时长(秒),比 3.1.1 的布尔值更灵活 | 第2章 |
| Keep Alive | Keep Alive(心跳间隔) | 客户端每隔多少秒发一个 PINGREQ 告诉 Broker “我还活着” | 第2章 |
| PINGREQ/PINGRESP | PING Request / Response | 心跳包:客户端发 PINGREQ,Broker 回 PINGRESP | 第2章 |
| PUB/SUB | Publish / Subscribe 报文 | MQTT 的两种核心报文:发消息 / 订阅主题 | 第2章 |
| 通配符 | Wildcard | + 匹配单层,# 匹配多层(订阅时用,发布时不能用) |
第2章 |
| ClientId | Client Identifier | 客户端的唯一标识,Broker 靠它认人(重复会导致互踢) | 第2章 |
| CONNECT/CONNACK | Connect / Connect Ack | 建立连接的报文对(客户端发 CONNECT,Broker 回 CONNACK) | 第2章 |
| PUBACK/PUBREC/PUBREL/PUBCOMP | QoS 确认报文 | QoS1 用 PUBACK;QoS2 用 PUBREC→PUBREL→PUBCOMP 四次握手 | 第2章 |
| EMQX | EMQX | 国产开源 MQTT Broker,性能好、中文文档全,国内首选 | 第2章 |
| Mosquitto | Eclipse Mosquitto | 老牌轻量 MQTT Broker,适合学习和小规模 | 第2章 |
| Paho | Eclipse Paho | Eclipse 出的 MQTT 客户端库(Java / Python / JS 等都有) | 第2章 |
| TLS | Transport Layer Security | 传输层加密,MQTT 走 TLS 后端口是 8883 | 第11章 |
| ACL | Access Control List | 权限表:控制哪个客户端能订阅/发布哪个主题 | 第11章 |
| OTA | Over-The-Air | 空中升级:通过 MQTT 远程给设备推送固件 | 第2章 |
| IoT | Internet of Things | 物联网,把设备连上网 | 第2章 |
| CoAP | Constrained Application Protocol | 另一个物联网协议,UDP 之上,比 MQTT 更轻但功能少 | 第2章 |
| 桥接 | Bridge | 两个 MQTT Broker 之间互相转发消息(跨集群/跨云) | 第2章 |
| 共享订阅 | Shared Subscription | 多个订阅者分摊同一个主题的消息(负载均衡,EMQX 扩展) | 第2章 |
目录
- 第一章:MQTT 到底是什么(小白必读)
- 第二章:MQTT 核心机制详解
- 第三章:代码实战(Broker 搭建 + 四种客户端)
- 第四章:应用场景、生产问题与解决方案、安全与性能
- 第五章:选型对比、面试题与小结
第一章:MQTT 到底是什么(小白必读)
这一章不写代码。搞懂“MQTT 是为谁设计的、和 HTTP 有什么区别”,后面所有机制都顺理成章。
1.1 从一个真实困境说起:HTTP 在物联网里哪里不行
假设你要做一个共享单车系统,全市 10 万辆车的锁要联网。每把锁每隔 30 秒上报一次位置和电量。
如果用 HTTP 轮询(锁主动请求服务器):
| 问题 | 具体表现 |
|---|---|
| 流量浪费 | 每次 HTTP 请求要带完整的 Header(Cookie、User-Agent…),可能 500 字节的头部只为传 20 字节的位置数据 |
| 耗电 | 锁是电池供电,HTTP 每次要建 TCP + TLS 握手,功耗是 MQTT 的 10 倍以上 |
| 服务器压力大 | 10 万设备每 30 秒一次 = 每秒 3300 次请求,还得维持 10 万个 TCP 连接 |
| 弱网卡顿 | 地下车库、隧道里信号差,HTTP 请求容易超时失败 |
| 设备无法被主动叫 | 服务器想远程开锁,HTTP 做不到(除非设备一直在轮询,更耗电) |
换成 MQTT 之后:
| 改善 | 原因 |
|---|---|
| 流量极小 | MQTT 最小报文只有 2 字节,头部极紧凑 |
| 省电 | 长连接(一次握手,之后一直用),不用反复建连 |
| 支持海量连接 | 单台 EMQX 可支撑 百万级并发连接 |
| 抗弱网 | 心跳机制 + 断线自动重连 + 离线消息(QoS1/2) |
| 双向通信 | 服务器能随时主动推消息给设备(远程开锁就靠这个) |
1.2 一句话定义 + 生活类比
定义:MQTT 是一个基于 TCP/IP 的、发布/订阅模式的、为低带宽 / 弱网 / 低功耗设备设计的超轻量消息协议。现在是 OASIS 国际标准(MQTT 3.1.1 和 MQTT 5.0)。
★ 生活类比:MQTT = 微信公众号 / 报社订阅
| MQTT 概念 | 公众号类比 | 说明 |
|---|---|---|
| Broker | 微信服务器 / 报社 | 中间机构,负责转发 |
| Publisher 发布者 | 公众号作者 | 写文章发出去 |
| Subscriber 订阅者 | 关注了这个号的用户 | 订阅后能收到 |
| Topic 主题 | 公众号名称 / 报纸栏目 | 你关注哪个号,就收哪个号的内容 |
| Publish 发布 | 作者发文 | 作者不知道谁会看到 |
| Subscribe 订阅 | 用户点关注 | 用户不知道作者是谁,只关心内容 |
| 保留消息 | 公众号的“置顶文章” | 新关注的人立刻能看到最新一篇 |
| 遗嘱消息 | 作者的“停更公告” | 作者账号突然注销,系统自动通知粉丝 |
★ 最关键的一点:发布者和订阅者互相不知道对方存在——这就是“解耦”。作者不用管有多少粉丝,粉丝不用管作者在哪。
1.3 MQTT 和 HTTP、RabbitMQ 的区别(面试必问)
这是最容易混淆的一组对比。先说结论:它们解决的是不同层面的问题。
| 维度 | MQTT | HTTP | RabbitMQ / Kafka |
|---|---|---|---|
| 通信模式 | 发布/订阅(Pub/Sub) | 请求/响应(Request/Response) | 发布/订阅 + 队列 |
| 连接方式 | 长连接(一次握手长期用) | 短连接(或 Keep-Alive 复用) | 长连接 |
| 谁主动 | 双向(服务端可主动推) | 客户端主动(服务端不能推) | 双向 |
| 报文头部 | 极小(最小 2 字节) | 大(几百字节) | 中等 |
| 设计目标 | 海量设备 + 弱网 + 省电 | 通用 Web | 后端服务解耦 + 高吞吐 |
| 消费模型 | 广播(所有订阅者都收到) | — | 队列=竞争消费;主题=广播 |
| 消息堆积 | ★ 能力弱(离线队列有限) | 无 | ★ 强项(能堆积海量消息) |
| 典型场景 | 物联网、移动推送、车联网 | 网页、API | 订单异步、削峰填谷 |
★ 最重要的三个区别(记这三个就够)
① MQTT 是“广播”,RabbitMQ 队列是“抢活”
【MQTT】 一个主题的消息,所有订阅者都会收到一份
设备 ──publish──► [topic/sensor] ──► 订阅者A ✓
──► 订阅者B ✓
──► 订阅者C ✓ (三人各收一份)
【RabbitMQ 队列】 一条消息,多个消费者"抢",只有一个拿到
生产者 ──► [queue.order] ──► 消费者A ✓
──► 消费者B ✗(没抢到)
──► 消费者C ✗(没抢到)
⚠️ 所以 MQTT 不能替代 RabbitMQ 做“任务分发”——因为 MQTT 会让所有消费者都干一遍同样的活。 (EMQX 有“共享订阅”扩展可以做到竞争消费,但那是扩展,不是标准 MQTT。)
② MQTT 消息不堆积,MQ 消息能堆积
MQTT 的设计假设是“设备在线就收,不在线就丢(或有限缓存)“。它不是为”堆积百万消息慢慢消费“设计的。 如果你的场景是”双十一订单积压几百万条慢慢处理“——那是 Kafka/RocketMQ 的活,别用 MQTT。
③ MQTT 服务端能主动推,HTTP 不能
这是物联网的核心需求:服务器要能随时给设备下发指令(开锁、升级、调参数)。 HTTP 只能设备主动问(轮询),MQTT 有长连接所以能直接推。
什么时候该用哪个(速查)
| 你的场景 | 用什么 |
|---|---|
| 设备上报传感器数据 | MQTT |
| 服务器远程控制设备(开锁/升级) | MQTT |
| App 消息推送(长连接) | MQTT(或专用推送通道) |
| 网页调后端 API | HTTP |
| 订单异步处理、削峰 | RabbitMQ / RocketMQ |
| 日志采集、大数据流 | Kafka |
| 设备数据进后端后的异步处理 | ★ MQTT 收数据 → 转 Kafka/RocketMQ 处理(常见架构!) |
★ 典型物联网架构:
设备 ──MQTT──► EMQX Broker ──(规则引擎)──► Kafka ──► 后端服务 ──► 数据库 (海量连接) (转接桥) (堆积/削峰)MQTT 负责“连得上”,MQ 负责“处理得过来”。两者是配合关系,不是替代关系。
1.4 发布/订阅模型详解
┌─────────────────────────────────────┐
│ MQTT Broker │
│ (EMQX / Mosquitto / HiveMQ) │
│ │
│ 主题树(Topic Tree) │
│ home/ │
│ ├─ livingroom/temp │
│ ├─ bedroom/temp │
│ └─ kitchen/humidity │
└─────────────────────────────────────┘
▲ publish │ subscribe
│ ▼
┌────────┴────────┐ ┌──────────────┐
│ 温度传感器 │ │ 手机 App │
│ (Publisher) │ │ (Subscriber) │
└─────────────────┘ └──────────────┘
发送 home/livingroom/temp = 26.5
★ 传感器完全不知道 App 的存在
关键点:
- 发布者只管往主题发,不关心谁订阅。
- 订阅者只管订阅主题,不关心谁发的。
- Broker 负责匹配和转发。
- 一个客户端既可以是发布者也可以是订阅者(比如一个智能灯:订阅“开灯指令”,发布“当前状态”)。
1.5 六个核心概念(先混个脸熟,第二章细讲)
| 概念 | 一句话 | 类比 |
|---|---|---|
| Topic 主题 | 消息的频道名,用 / 分层 |
报纸栏目:“体育/足球/中超” |
| QoS 服务质量 | 消息送达的可靠等级(0/1/2) | 寄信:平信 / 挂号信 / 回执挂号信 |
| Retained 保留消息 | Broker 记住每个主题的最后一条 | 报亭摆着最新一期,新来的人立刻拿到 |
| LWT 遗嘱消息 | 设备掉线时 Broker 代发的“遗言” | 登山前的遗书,出事了自动公开 |
| Session 会话 | Broker 为客户端保留的档案 | 邮局的“留局待取”服务 |
| Keep Alive 心跳 | 定期告诉 Broker 我还活着 | 每隔一段时间打个电话报平安 |
1.6 MQTT 版本:3.1.1 还是 5.0?
| 版本 | 年份 | 状态 | 说明 |
|---|---|---|---|
| MQTT 3.1 | 2010 | 过时 | 最初版本 |
| MQTT 3.1.1 | 2014 | ★ 最广泛支持 | 目前兼容性最好,绝大多数 Broker/客户端都支持 |
| MQTT 5.0 | 2019 | 新特性多 | 增加了原因码、会话过期、共享订阅、消息过期、用户属性等 |
MQTT 5.0 的主要增强(面试能说出来加分):
| 特性 | 作用 |
|---|---|
| Reason Code 原因码 | 所有响应都带原因码,出错能知道为什么(3.1.1 只有一个笼统的错误码) |
| Session Expiry Interval | 会话可以设过期时间(秒),比 3.1.1 的 cleanSession 布尔值灵活 |
| Message Expiry Interval | 消息可以设过期时间,过期的离线消息不再投递 |
| Shared Subscription | 原生支持共享订阅(负载均衡式消费) |
| User Properties | 消息可以带自定义 KV 头(像 HTTP Header) |
| Topic Alias | 主题别名,用数字代替长主题名,省流量 |
| Flow Control | 流量控制,防止快客户端把慢 Broker 压垮 |
选型建议:
- 新项目、Broker 和客户端都支持 → MQTT 5.0
- 设备端 SDK 老旧 / 要兼容老设备 → MQTT 3.1.1
- 国内现状:EMQX 同时支持两者,设备端用 3.1.1 居多(SDK 兼容性),服务端应用侧可用 5.0
1.7 什么时候该用 MQTT,什么时候不该用
✅ 适合
| 场景 | 为什么 |
|---|---|
| 设备上报数据(传感器、电表、车联网) | 海量连接、小报文、省电 |
| 远程控制设备(开锁、调参数、OTA 升级) | 服务端要能主动推 |
| 弱网环境(地下车库、野外、海上) | 心跳 + 断线重连 + 离线消息 |
| App 长连接推送(IM、直播弹幕) | 长连接双向 |
| 移动设备(省流量省电) | 头部极小 |
❌ 不适合
| 场景 | 为什么 |
|---|---|
| 需要消息堆积削峰 | MQTT 不擅长堆积,用 Kafka/RocketMQ |
| 需要严格的事务消息 | MQTT 没有事务概念,用 RocketMQ |
| 需要竞争消费(一条消息只被处理一次) | MQTT 是广播;用 RabbitMQ 队列(或 EMQX 共享订阅) |
| 后端微服务之间调用 | 用 HTTP/gRPC/RPC 框架,MQTT 是设备侧协议 |
| 大数据量传输(传文件、视频) | MQTT 适合小报文(默认限制 256MB 但实践上应 <1MB) |
1.8 本章小结
┌──────────────────────────────────────────────────────────┐
│ MQTT 一句话总结 │
│ │
│ "为弱网、低功耗、海量设备设计的、发布订阅模式的、 │
│ 超轻量长连接消息协议" │
│ │
│ 和 HTTP 比:长连接 + 服务端可推送 + 头部极小 │
│ 和 MQ 比:广播而非竞争、不擅长堆积、面向设备而非服务 │
│ │
│ 典型架构:设备 ──MQTT──► Broker ──► Kafka ──► 业务系统 │
└──────────────────────────────────────────────────────────┘
第二章:MQTT 核心机制详解
这一章是 MQTT 的灵魂。QoS、主题、保留消息、遗嘱、会话——这五个搞懂,你就超过 80% 的候选人。
2.1 报文结构:为什么它这么轻
MQTT 每个报文都由三部分组成:
┌──────────────┬──────────────┬──────────────┐
│ 固定头 │ 可变头 │ 有效载荷 │
│ Fixed Header │ Variable Hdr │ Payload │
│ (2~5 字节) │ (部分报文有) │ (部分报文有) │
└──────────────┴──────────────┴──────────────┘
固定头(每个报文必有,最小 2 字节):
Bit: 7 6 5 4 3 2 1 0
┌─────────┬─────────────┐
│报文类型 │ 标志位 Flags │ ← 第 1 字节
├─────────┴─────────────┤
│ 剩余长度 Remaining Length │ ← 第 2~5 字节(变长编码,1~4 字节)
└───────────────────────┘
★ 为什么最小只有 2 字节:像 PINGREQ、DISCONNECT 这种报文,只有固定头、没有可变头和载荷,所以就是 2 字节。对比 HTTP 动辄几百字节的头部,省了几十倍。
14 种报文类型(面试能说出主要的就行)
| 值 | 报文 | 方向 | 作用 |
|---|---|---|---|
| 1 | CONNECT | 客户端 → Broker | 请求建立连接 |
| 2 | CONNACK | Broker → 客户端 | 连接确认(返回连接结果码) |
| 3 | PUBLISH | 双向 | 发布消息(唯一真正传数据的报文) |
| 4 | PUBACK | 双向 | QoS 1 的确认 |
| 5 | PUBREC | 双向 | QoS 2 第一步:已收到(Received) |
| 6 | PUBREL | 双向 | QoS 2 第二步:已释放(Release) |
| 7 | PUBCOMP | 双向 | QoS 2 第三步:已完成(Complete) |
| 8 | SUBSCRIBE | 客户端 → Broker | 订阅主题 |
| 9 | SUBACK | Broker → 客户端 | 订阅确认(返回授予的 QoS) |
| 10 | UNSUBSCRIBE | 客户端 → Broker | 取消订阅 |
| 11 | UNSUBACK | Broker → 客户端 | 取消订阅确认 |
| 12 | PINGREQ | 客户端 → Broker | 心跳请求 |
| 13 | PINGRESP | Broker → 客户端 | 心跳响应 |
| 14 | DISCONNECT | 客户端 → Broker | 断开连接 |
| 15 | AUTH(5.0) | 双向 | 增强认证(MQTT 5.0 新增) |
💡 记忆:
连接 1-2、发布 3-7(3 是数据、4 是 QoS1、567 是 QoS2)、订阅 8-11、心跳 12-13、断开 14。
2.2 QoS 三级服务质量(★ 面试必考,也最容易用错)
QoS = Quality of Service,消息投递的可靠等级。三个级别:
| QoS | 中文 | 英文 | 交付保证 | 报文开销 | 类比 |
|---|---|---|---|---|---|
| 0 | 最多一次 | At most once | 发出去就不管,可能丢 | 1 个 PUBLISH | 平信:扔进邮筒,丢不丢看运气 |
| 1 | 至少一次 | At least once | 保证送到,但可能重复 | PUBLISH + PUBACK | 挂号信:要签收回执,没收到就重发(可能寄重复了) |
| 2 | 恰好一次 | Exactly once | 保证送到且不重复 | 4 次握手 | 回执挂号信:来回确认四次,确保只收到一份 |
QoS 0:发完就忘(Fire and Forget)
Publisher ────PUBLISH(QoS0)────► Broker ────PUBLISH(QoS0)────► Subscriber
(发完不管,没有确认)
特点:最快、最省流量、但可能丢消息(网络断了就丢了)。
什么时候用:传感器高频上报温度——丢一两条无所谓,反正下一条马上来。
QoS 1:至少一次(有确认,可能重复)
Publisher Broker
│ │
│───────── PUBLISH(QoS1, msgId=1) ──────────►│ ① 发出并暂存
│ │
│◄──────────── PUBACK(msgId=1) ──────────────│ ② 确认收到
│ │
│ (超时没收到 PUBACK → 重发,msgId 相同, │
│ 但 DUP 标志位置 1) │
│───────── PUBLISH(QoS1, msgId=1, DUP=1) ───►│ ③ 重发
│◄──────────── PUBACK(msgId=1) ──────────────│
★ 关键理解:“至少一次”意味着消息一定到达,但可能到达多次。
为什么可能重复:
- 发送方发出 PUBLISH 后,在收到 PUBACK 之前连接断了。
- 发送方重连后,不确定 Broker 到底收到没有,只能重发(DUP=1)。
- 如果 Broker 其实已经收到了,就会投递两次给订阅者。
⚠️ 所以 QoS 1 的消费端必须做幂等! 这是生产上最常见的坑(见 4.2 问题 2)。
什么时候用:设备状态上报、告警——不能丢,重复了业务能容忍(或做幂等处理)。
QoS 2:恰好一次(四次握手)
Publisher Broker
│ │
│─────────── PUBLISH(QoS2, msgId=1) ───────────►│ ① 发出,暂存消息
│ │
│◄────────────── PUBREC(msgId=1) ──────────────│ ② "我收到了"(Received)
│ │
│────────────── PUBREL(msgId=1) ───────────────►│ ③ "你可以释放了"(Release)
│ │
│◄────────────── PUBCOMP(msgId=1) ─────────────│ ④ "已完成"(Complete)
│ │
│ 双方都删掉暂存的 msgId=1 │
★ 为什么四次握手能保证“恰好一次”:
- 发送方用
PUBREC记住“这个消息 Broker 已收到但还没处理完”,之后即使重发也用同一个 msgId。 PUBREL之后发送方确认可以丢弃暂存了。- Broker 在
PUBREL之前不会投递给订阅者,PUBREL之后只投递一次。
特点:最可靠,但开销是 QoS 1 的 2 倍(4 个报文)、延迟更高。
什么时候用:计费、支付指令、关键控制指令——绝对不能丢也不能重复。
★ QoS 降级规则(面试高频陷阱)
⚠️ ⚠️ QoS 是“发布 QoS”和“订阅 QoS”取小值!
发布者用 QoS 2 发布 ──► Broker ──► 订阅者用 QoS 0 订阅
↓
订阅者实际收到的是 QoS 0(降级了!)
| 发布 QoS | 订阅 QoS | 实际投递 QoS |
|---|---|---|
| 2 | 2 | 2 |
| 2 | 1 | 1(降级) |
| 2 | 0 | 0(降级) |
| 1 | 2 | 1(取小的) |
| 0 | 2 | 0(降级) |
★ 面试话术:“QoS 有两个端——发布者到 Broker 是一段,Broker 到订阅者是另一段,两段各自独立协商,最终投递给订阅者的是两者取小值。 所以想保证端到端 QoS 2,发布和订阅都必须设成 2。很多线上’消息丢了’的事故,就是因为发布端设了 QoS 1,但订阅端用了默认的 QoS 0。”
QoS 选型速查
| 场景 | 推荐 QoS | 理由 |
|---|---|---|
| 高频传感器数据(温度、湿度) | 0 | 丢了无所谓,下一条马上来 |
| 设备状态、门禁刷卡记录 | 1 | 不能丢,重复可幂等 |
| 告警(火灾、入侵) | 1 或 2 | 不能丢 |
| 计费、支付、交易指令 | 2 | 不能丢也不能重复 |
| 远程开锁 / 关阀 | 1(+ 应用层确认) | 要可靠,且业务上有“执行结果回报”兜底 |
⚠️ 诚实提醒:即使 QoS 2,也只能保证消息层面的“恰好一次”,不能保证业务层面的“恰好执行一次”。 比如设备收到“开锁”指令后开完锁就宕机了,没回报执行结果——你不知道锁开没开。 真正的业务可靠性要靠应用层设计(指令 ID + 执行结果回报 + 状态查询)。
2.3 主题(Topic)与通配符
主题的层级结构
Topic 是用 / 分隔的字符串(UTF-8),像文件路径:
home/livingroom/temperature
company/building1/floor3/room301/light
vehicle/VIN123456/gps
device/D001/status
设计建议:
| 原则 | 示例 | 说明 |
|---|---|---|
不要以 / 开头 |
✅ home/temp ❌ /home/temp |
开头斜杠会产生一个空层级 |
| 不要有空格 | ✅ living_room ❌ living room |
空格容易出解析问题 |
| 区分大小写 | Home/Temp ≠ home/temp |
建议全小写统一 |
| 别太长 | 建议 < 128 字节 | 每条消息都带主题名,长了费流量 |
| 把不变的部分放前面 | device/{id}/status |
便于通配符订阅和 ACL 控制 |
两个通配符(★ 只在订阅时用)
| 通配符 | 含义 | 示例 | 能匹配 |
|---|---|---|---|
+ |
单层通配 | home/+/temp |
home/livingroom/temp、home/bedroom/temp❌ 不匹配 home/livingroom/sensor1/temp |
# |
多层通配(必须在最后) | home/# |
home/temp、home/livingroom/temp、home/a/b/c |
// 订阅示例
client.subscribe("home/+/temp", 1); // 订阅所有房间的温度
client.subscribe("device/#", 1); // 订阅 device 下所有主题
client.subscribe("vehicle/+/gps", 1); // 订阅所有车辆的 GPS
client.subscribe("#", 0); // ⚠️ 订阅所有消息(危险!调试时才用)
⚠️ 坑 1:发布时不能用通配符。你必须发布到具体的主题(
home/livingroom/temp),publish("home/+/temp", ...)会发布到一个字面上叫home/+/temp的主题,没有订阅者能收到。⚠️ 坑 2:通配符订阅会给 Broker 带来额外开销。
#订阅在多主题场景下性能很差,生产上慎用。
特殊主题 $
以 $ 开头的主题是 Broker 的系统主题,通常不会被 # 通配符匹配到:
$SYS/brokers/emqx@127.0.0.1/version Broker 版本
$SYS/brokers/emqx@127.0.0.1/uptime 运行时长
$SYS/brokers/emqx@127.0.0.1/clients/count 当前连接数
$share/group1/device/# 共享订阅(见下)
⚠️ 坑:订阅
#收不到$SYS/...的消息——这是协议规定的(防止系统主题被误订阅)。要监控 Broker 得显式订阅$SYS/#。
共享订阅(负载均衡式消费)
标准 MQTT 是广播(所有订阅者都收到)。共享订阅让多个订阅者分摊消息:
$share/group1/device/#
│
├─► 订阅者 A ← 消息 1、4、7...
├─► 订阅者 B ← 消息 2、5、8...
└─► 订阅者 C ← 消息 3、6、9...
// 三个消费者订阅同一个共享组,消息会被分摊
client.subscribe("$share/orderGroup/device/data", 1);
💡 面试加分:标准 MQTT 没有共享订阅,这是 MQTT 5.0 特性 / EMQX 扩展。它让 MQTT 也能做“竞争消费”,是 MQTT 桥接后端 MQ 之外的一种解法。
2.4 保留消息(Retained Message)
一句话:Broker 会记住每个主题的最后一条消息(带 retain 标志的),新订阅者一上线就立刻收到,不用等下次发布。
★ 生活类比:报亭的“最新一期样刊”——你第一次去这家报亭,老板直接把最新一期给你,你不用等到明天发新刊。
【没有保留消息】
① 温度传感器发布 home/temp = 26.5(retain=false)
② Broker 转发给当前订阅者,然后就忘了
③ 手机 App 打开,订阅 home/temp
④ App 什么都收不到,要等传感器下一次上报(可能 30 秒后)
→ 用户体验:打开 App 一片空白,干等
【有保留消息】
① 温度传感器发布 home/temp = 26.5(retain=true)
② Broker 转发,并且**记住**这条
③ 手机 App 打开,订阅 home/temp
④ Broker 立刻把 26.5 推给 App ✓
→ 用户体验:打开 App 立刻看到当前温度
代码:
// 发布保留消息(第 3 个参数 qos,第 4 个参数 retained=true)
client.publish("home/livingroom/temp", "26.5".getBytes(), 1, true);
// 清除某个主题的保留消息:发一条空消息 + retained=true
client.publish("home/livingroom/temp", new byte[0], 1, true);
⚠️ 坑:保留消息每个主题只保留一条,新发布的会覆盖旧的。 要删除就发空载荷 + retained=true(不能发 null,必须是长度为 0 的字节数组)。
什么时候用:设备最新状态(灯是开是关、温度多少、在线状态)——新上线的 App 要立刻看到当前状态。
2.5 遗嘱消息 LWT(Last Will and Testament)
一句话:设备连接时提前登记一条“遗言”,一旦它非正常掉线,Broker 就自动把这条消息发出去。
★ 生活类比:登山前给家人留一封信——“如果我 48 小时没联系你们,就把这封信公开”。
【场景】10 万个设备,某个设备突然断电/进电梯没信号
设备连接时:CONNECT 里带
willTopic = "device/D001/status"
willMessage = "{\"online\":false}"
willQoS = 1
willRetain = true
│
▼
Broker 记住这份遗嘱(但不发)
│
▼
设备突然断电(没发 DISCONNECT,TCP 断)
│
▼
Broker 检测到连接异常断开 → 自动发布遗嘱消息到 device/D001/status
│
▼
监控系统的订阅者立刻收到 {"online":false} → 标记设备离线 → 告警
代码:
MqttConnectOptions opts = new MqttConnectOptions();
opts.setUserName("device001");
opts.setPassword("secret".toCharArray());
// ★ 设置遗嘱
opts.setWill("device/D001/status", // 遗嘱主题
"{\"online\":false}".getBytes(), // 遗嘱内容
1, // QoS
true); // retained(保留,新订阅者也能看到"离线"状态)
// 心跳间隔(秒)—— Broker 多久没收到消息就判定掉线
opts.setKeepAliveInterval(60);
client.connect(opts);
// ★ 设备正常启动时,主动发布"上线"(覆盖遗嘱)
client.publish("device/D001/status", "{\"online\":true}".getBytes(), 1, true);
★ 关键设计:设备上线后要主动发一条“上线”消息(retained=true),覆盖掉遗嘱。这样:
正常流程:
连接(带遗嘱 offline) → 发 retained "online" → ... → 正常断开(发 DISCONNECT)
↓
主动发 retained "offline" 再断开
(这样不会触发遗嘱,但状态已正确)
异常流程:
连接(带遗嘱 offline) → 发 retained "online" → 突然断电
↓
Broker 检测心跳超时 → 自动发遗嘱 "offline" ✓
⚠️ 坑 1(非常重要):正常调用
disconnect()不会触发遗嘱! 遗嘱只在异常断开时触发(TCP 断了、心跳超时)。所以正常下线时要自己发一条 offline 消息,否则状态会一直是 online。⚠️ 坑 2:遗嘱是在 CONNECT 报文里设置的,连接建立后不能修改。要改遗嘱内容,只能断开重连。
⚠️ 坑 3:如果 Broker 也崩了,遗嘱不会发——遗嘱是 Broker 代发的,Broker 自己挂了就没人发了。
2.6 会话(Session):断线重连后还能收到消息吗
会话 = Broker 为客户端保留的“档案”,包含:
- 客户端的订阅关系
- 已发送但未确认的消息(QoS 1/2)
- 已收到但未完成的消息(QoS 2)
MQTT 3.1.1:Clean Session 布尔值
MqttConnectOptions opts = new MqttConnectOptions();
opts.setCleanSession(true); // ① true:不保留会话,断开即忘
opts.setCleanSession(false); // ② false:保留会话,重连后继续
| CleanSession | 断开后 | 重连后 | 适用场景 |
|---|---|---|---|
| true | Broker 删除会话、订阅关系、未确认消息 | 要重新订阅,离线期间的消息丢失 | 临时连接、只发不收、设备资源紧张 |
| false | Broker 保留会话 | 订阅关系自动恢复,离线期间的 QoS1/2 消息补发 | ★ 设备需要“离线消息” |
★ 离线消息是怎么工作的:
设备(CleanSession=false, ClientId=D001)订阅 device/D001/cmd, QoS=1
│
▼
设备断线(进电梯)
│
▼
服务器发布 3 条指令到 device/D001/cmd(QoS=1)
│
▼
Broker 发现 D001 不在线,但有会话 → 把消息**暂存在会话里**
│
▼
设备重连(同样的 ClientId=D001,CleanSession=false)
│
▼
Broker 恢复会话 → 自动恢复订阅 → 把暂存的 3 条指令**补发**给设备 ✓
⚠️ 坑:离线消息不是无限的。Broker 会限制队列长度(EMQX 默认几万条),超了会丢弃最老的。 而且 ★ 只有 QoS 1/2 的消息才会被暂存,QoS 0 的离线消息直接丢弃。
MQTT 5.0:Session Expiry Interval(更灵活)
// 5.0:会话保留 1 小时(3600 秒),之后自动清理
opts.setSessionExpiryInterval(3600L);
// 0 = 断开即清理(等价于 CleanSession=true)
// 0xFFFFFFFF = 永不过期(等价于 CleanSession=false)
比 3.1.1 好在哪:3.1.1 里 CleanSession=false 会让会话永久保留,设备刷机后再也不上线,Broker 上就堆了一堆僵尸会话(内存泄漏风险)。
5.0 可以设“保留 1 天”,过期自动清理——这是运维上的重要改进。
★ ClientId 的坑(生产事故高发)
Broker 靠 ClientId 识别客户端。两个连接用同一个 ClientId,后连的会把先连的踢掉。
设备 A (ClientId=D001) 已连接
↓
设备 B (ClientId=D001) 也来连接
↓
Broker: "同一个 ID,只能有一个" → 断开 A,接受 B
↓
设备 A 收到"被踢"→ 重连 → 又踢掉 B
↓
★ 两个设备互相踢,无限循环 —— 经典"互踢"故障!
解决方案:
// ❌ 错误:所有设备用同一个 ClientId(抄 Demo 抄出来的)
opts.setClientId("mqtt_client");
// ✅ 正确:用设备唯一标识 + 后缀
opts.setClientId("D001_" + deviceSn); // 设备序列号
opts.setClientId("D001_" + macAddress.replace(":", "")); // MAC 地址
opts.setClientId("D001_" + UUID.randomUUID()); // 不要持久化会话时才用随机
⚠️ 注意:如果你需要离线消息(CleanSession=false),ClientId 必须固定(不能随机)——因为 Broker 靠 ClientId 找回会话。 所以:
ClientId 固定 + 保证每设备唯一是唯一正确解。
2.7 心跳(Keep Alive)与 PINGREQ
机制:
- 客户端 CONNECT 时声明
KeepAlive秒数(如 60)。 - 客户端必须在 1.5 × KeepAlive 时间内至少发一个报文(可以是 PINGREQ,也可以是任意 PUBLISH)。
- 如果客户端没数据要发,就发 PINGREQ,Broker 回 PINGRESP。
- Broker 超时没收到任何报文 → 判定客户端掉线 → 触发遗嘱 + 关闭连接。
客户端 Broker
│ │
│────── CONNECT(KeepAlive=60) ──────────►│
│◄────────── CONNACK ────────────────────│
│ │
│ (60 秒内没有数据要发) │
│────────── PINGREQ ────────────────────►│ ← 我还在
│◄───────── PINGRESP ───────────────────│ ← 好的
│ │
│ (又安静了 60 秒) │
│────────── PINGREQ ────────────────────►│
│◄───────── PINGRESP ───────────────────│
Keep Alive 怎么设:
| 场景 | 建议值 | 理由 |
|---|---|---|
| 稳定网络、服务端应用 | 60 秒 | 常规值 |
| 移动网络 / 弱网 | 30~60 秒 | 短一点能更快发现掉线 |
| 电池供电设备 | 300~1800 秒 | ★ 心跳越频繁越耗电! |
| 需要快速感知掉线 | 10~30 秒 | 车联网、实时监控 |
⚠️ 权衡:心跳短 = 掉线发现快,但耗电 + 费流量。 电池设备(如烟感报警器)可能设 30 分钟心跳,代价是掉线后最多 45 分钟才被发现。
⚠️ 坑:Paho Java 客户端的心跳是客户端库自动发的,你不用手写。但如果你的业务线程阻塞了(比如同步等一个慢接口),心跳发不出去,会被 Broker 判定掉线——所以 MQTT 回调里不能做耗时操作(见 4.2 问题 9)。
2.8 连接建立与断开的完整流程
【建立连接】
客户端 Broker
│ │
│──────── TCP 三次握手 ───────────────────────►│
│──────── CONNECT ───────────────────────────►│
│ ├ clientId │
│ ├ username / password(可选) │ ① Broker 认证
│ ├ cleanSession / sessionExpiry │ ② 检查 ClientId 唯一性
│ ├ keepAlive │ ③ 建立会话
│ └ will(遗嘱,可选) │
│◄─────── CONNACK ────────────────────────────│
│ ├ returnCode(0=成功,非0看原因) │
│ └ sessionPresent(会话是否恢复) │
│ │
│──────── SUBSCRIBE(topic, qos) ─────────────►│
│◄─────── SUBACK(grantedQoS) ─────────────────│
│ │
│══════════ 开始收发消息 ══════════════════════│
【正常断开】
│──────── DISCONNECT ────────────────────────►│ ★ 不会触发遗嘱
│──────── TCP 关闭 ──────────────────────────►│
CONNACK 返回码(3.1.1):
| 码 | 含义 | 排查方向 |
|---|---|---|
| 0 | 连接成功 | — |
| 1 | 不支持的协议版本 | 客户端/服务端 MQTT 版本不匹配 |
| 2 | ClientId 被拒绝 | ClientId 为空或格式非法 |
| 3 | 服务端不可用 | Broker 故障 |
| 4 | 用户名或密码错误 | ★ 最常见:账号密码错 / 没配认证 |
| 5 | 未授权 | ★ 认证通过但 ACL 不允许连接 |
💡 MQTT 5.0 的改进:5.0 的 CONNACK 带 Reason Code + 可选的原因字符串,比如“密码错误”、“ClientId 重复”,排查问题比 3.1.1 的笼统码方便得多。
第三章:代码实战(Broker 搭建 + 四种客户端)
这一章的代码全部可以直接跑。建议先按 3.1 起一个 Broker,然后边读边跑。
3.1 搭建 MQTT Broker
3.1.1 用 Docker 起 EMQX(推荐)
# EMQX 5.x(国产,性能好,自带 Web 控制台,中文文档全)
docker run -d --name emqx \
-p 1883:1883 \ # MQTT TCP 端口
-p 8883:8883 \ # MQTT over TLS
-p 8083:8083 \ # MQTT over WebSocket
-p 8084:8084 \ # MQTT over WSS
-p 18083:18083 \ # Web 管理控制台
emqx/emqx:5.7.0
# 打开控制台:http://localhost:18083
# 默认账号:admin / public(首次登录会强制改密码)
3.1.2 用 Docker 起 Mosquitto(轻量,学习用)
docker run -it --name mosquitto \
-p 1883:1883 \
-p 9001:9001 \
eclipse-mosquitto:2.0
3.1.3 用命令行快速测试(不用写代码)
# 订阅(终端 1)
mosquitto_sub -h localhost -p 1883 -t "home/#" -v
# 发布(终端 2)
mosquitto_pub -h localhost -p 1883 -t "home/livingroom/temp" -m "26.5"
# 带用户名密码
mosquitto_sub -h localhost -t "test" -u "admin" -P "public"
# 指定 QoS = 1
mosquitto_pub -h localhost -t "test" -m "hello" -q 1
# 发布保留消息
mosquitto_pub -h localhost -t "home/status" -m "online" -r
3.1.4 端口速查
| 端口 | 用途 |
|---|---|
| 1883 | MQTT over TCP(明文) |
| 8883 | MQTT over TLS(加密)★ 生产必用 |
| 8083 | MQTT over WebSocket(浏览器用) |
| 8084 | MQTT over WSS |
| 18083 | EMQX Web 控制台 |
3.2 Java 客户端:Eclipse Paho
3.2.1 Maven 依赖
<dependency>
<groupId>org.eclipse.paho</groupId>
<artifactId>org.eclipse.paho.client.mqttv3</artifactId>
<version>1.2.5</version>
</dependency>
3.2.2 完整客户端(★ 生产级:含自动重连、回调、离线缓冲)
/**
* 生产级 MQTT 客户端:自动重连 + 回调处理 + 遗嘱 + 离线缓冲
*/
@Component
@Slf4j
public class MqttClientService {
private MqttClient client;
private final String broker = "tcp://localhost:1883";
private final String clientId = "app-server-001";
@PostConstruct
public void init() {
try {
client = new MqttClient(broker, clientId, new MqttClientPersistence() {
// 用内存持久化(生产可用文件/Redis 持久化,见下方说明)
private final Map<String, MqttPersistable> store = new ConcurrentHashMap<>();
@Override public void open(String s, String s1) {}
@Override public void close() {}
@Override public void put(String key, MqttPersistable p) { store.put(key, p); }
@Override public MqttPersistable get(String key) { return store.get(key); }
@Override public void remove(String key) { store.remove(key); }
@Override public Enumeration<String> keys() { return Collections.enumeration(store.keySet()); }
@Override public void clear() { store.clear(); }
@Override public boolean containsKey(String key) { return store.containsKey(key); }
});
// ★ 设置回调(连接丢失、收到消息、发送完成)
client.setCallback(new MqttCallbackExtended() {
/**
* 连接完成(包括自动重连成功)★ 重点
* reconnect=true 表示这是自动重连,需要在这里重新订阅
*/
@Override
public void connectComplete(boolean reconnect, String serverURI) {
log.info("MQTT 连接成功,serverURI={}, 是否重连={}", serverURI, reconnect);
try {
// ★ 重连后必须重新订阅(除非用了 cleanSession=false 且 Broker 保留了会话)
client.subscribe("device/+/data", 1);
client.subscribe("device/+/status", 1);
log.info("订阅完成");
} catch (MqttException e) {
log.error("重连后订阅失败", e);
}
}
/** 连接丢失 —— Paho 会自动重连,这里只做日志 */
@Override
public void connectionLost(Throwable cause) {
log.warn("MQTT 连接丢失,Paho 将自动重连", cause);
}
/** 收到消息 ★ 核心业务处理在这 */
@Override
public void messageArrived(String topic, MqttMessage message) {
String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
log.info("收到消息 topic={}, qos={}, payload={}",
topic, message.getQos(), payload);
try {
handleMessage(topic, payload);
} catch (Exception e) {
// ⚠️ 必须 catch,否则异常抛到 Paho 线程会导致连接断开
log.error("处理 MQTT 消息失败, topic={}", topic, e);
}
}
/** QoS 1/2 消息发送完成(收到 PUBACK/PUBCOMP) */
@Override
public void deliveryComplete(IMqttDeliveryToken token) {
log.debug("消息投递完成, msgId={}", token.getMessageId());
}
});
MqttConnectOptions opts = buildConnectOptions();
client.connect(opts);
log.info("MQTT 客户端启动成功,clientId={}", clientId);
} catch (MqttException e) {
log.error("MQTT 连接失败", e);
}
}
private MqttConnectOptions buildConnectOptions() {
MqttConnectOptions opts = new MqttConnectOptions();
// ① 认证
opts.setUserName("admin");
opts.setPassword("public".toCharArray());
// ② ★ 自动重连(默认就是 true,显式写出来更清晰)
opts.setAutomaticReconnect(true);
// ③ ★ 会话:false = 保留会话,重连后补发离线消息
opts.setCleanSession(false);
// ④ 心跳(秒)
opts.setKeepAliveInterval(60);
// ⑤ 连接超时(秒)
opts.setConnectionTimeout(30);
// ⑥ ★ 遗嘱:本服务意外掉线时通知大家
opts.setWill("app/server/status",
"{\"online\":false}".getBytes(StandardCharsets.UTF_8),
1, true);
// ⑦ 离线缓冲:断线期间要发的消息先存起来,重连后自动发
// 注意:只有 QoS 1/2 且 cleanSession=false 才有意义
opts.setMaxInflight(1000); // 允许同时有多少条未确认消息在飞
return opts;
}
/** 发布消息 */
public void publish(String topic, String payload, int qos, boolean retained) {
try {
MqttMessage message = new MqttMessage(payload.getBytes(StandardCharsets.UTF_8));
message.setQos(qos);
message.setRetained(retained);
client.publish(topic, message);
} catch (MqttException e) {
log.error("发布失败 topic={}", topic, e);
}
}
/** ★ 下发指令给设备(QoS 1 + 应用层超时等待回报) */
public boolean sendCommand(String deviceId, String command, long timeoutMs) {
try {
String topic = "device/" + deviceId + "/cmd";
MqttMessage msg = new MqttMessage(command.getBytes(StandardCharsets.UTF_8));
msg.setQos(1); // 至少一次,保证送达
// 同步发送:等 PUBACK
IMqttDeliveryToken token = client.publishAndWait(topic, msg, timeoutMs);
return token.isComplete();
} catch (MqttException e) {
log.error("下发指令失败 deviceId={}", deviceId, e);
return false;
}
}
/** 业务处理 */
private void handleMessage(String topic, String payload) {
// 按主题分发
if (topic.endsWith("/data")) {
// 解析设备数据 → 入库 / 转发 Kafka
DeviceData data = JSON.parseObject(payload, DeviceData.class);
deviceService.save(data);
} else if (topic.endsWith("/status")) {
// 更新设备在线状态
deviceService.updateOnlineStatus(payload);
}
}
@PreDestroy
public void destroy() {
try {
// ★ 正常下线:先发 offline 状态,再断开(避免状态残留)
publish("app/server/status", "{\"online\":false}", 1, true);
if (client != null && client.isConnected()) {
client.disconnect();
client.close();
}
} catch (MqttException e) {
log.error("关闭 MQTT 连接失败", e);
}
}
}
3.2.3 关键 API 速查
| 方法 | 作用 |
|---|---|
new MqttClient(broker, clientId, persistence) |
创建客户端 |
setCallback(MqttCallbackExtended) |
设置回调(推荐 Extended 版,有 connectComplete) |
connect(MqttConnectOptions) |
连接 |
subscribe(topic, qos) |
订阅 |
publish(topic, MqttMessage) |
发布(异步) |
publishAndWait(...) |
发布并同步等待确认(要等 PUBACK) |
disconnect() |
正常断开(不触发遗嘱) |
isConnected() |
是否已连接 |
close() |
释放资源 |
3.2.4 TLS 加密连接(生产必配)
public MqttConnectOptions buildSslOptions() throws Exception {
MqttConnectOptions opts = new MqttConnectOptions();
opts.setUserName("admin");
opts.setPassword("public".toCharArray());
// ★ 端口用 8883,协议用 ssl://
// MqttClient client = new MqttClient("ssl://broker.example.com:8883", clientId);
// 配置 SSL
SSLContext sslContext = SSLContext.getInstance("TLSv1.2");
TrustManagerFactory tmf = TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm());
KeyStore ks = KeyStore.getInstance("JKS");
try (InputStream is = new FileInputStream("ca.jks")) {
ks.load(is, "changeit".toCharArray()); // 装载 CA 证书
}
tmf.init(ks);
sslContext.init(null, tmf.getTrustManagers(), null);
opts.setSocketFactory(sslContext.getSocketFactory());
return opts;
}
3.3 Spring Integration MQTT(注解驱动,推荐 Spring 项目用)
3.3.1 依赖与配置
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-mqtt</artifactId>
</dependency>
# application.yml
mqtt:
broker-url: tcp://localhost:1883
client-id: spring-app-${random.value} # ★ 避免多实例 ClientId 冲突
username: admin
password: public
default-topic: device/#
timeout: 30
keepalive: 60
3.3.2 配置类
@Configuration
@Slf4j
public class MqttConfig {
@Value("${mqtt.broker-url}") private String brokerUrl;
@Value("${mqtt.client-id}") private String clientId;
@Value("${mqtt.username}") private String username;
@Value("${mqtt.password}") private String password;
@Value("${mqtt.default-topic}") private String defaultTopic;
/** ① 连接工厂 */
@Bean
public MqttPahoClientFactory mqttClientFactory() {
DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
MqttConnectOptions options = new MqttConnectOptions();
options.setServerURIs(new String[]{brokerUrl});
options.setUserName(username);
options.setPassword(password.toCharArray());
options.setCleanSession(false); // 保留会话
options.setAutomaticReconnect(true); // 自动重连
options.setKeepAliveInterval(60);
options.setMaxInflight(1000);
factory.setConnectionOptions(options);
return factory;
}
/** ② 入站通道适配器(接收消息) */
@Bean
public MessageProducer inbound() {
// 订阅两个主题,QoS 都是 1
MqttPahoMessageDrivenChannelAdapter adapter =
new MqttPahoMessageDrivenChannelAdapter(
clientId + "-in", mqttClientFactory(),
"device/+/data", "device/+/status");
adapter.setCompletionTimeout(5000);
adapter.setConverter(new DefaultPahoMessageConverter());
adapter.setQos(1); // ★ 订阅 QoS
adapter.setOutputChannel(mqttInputChannel());
return adapter;
}
@Bean
public MessageChannel mqttInputChannel() {
return new DirectChannel();
}
/** ③ 消息处理器 ★ 业务代码写这 */
@Bean
@ServiceActivator(inputChannel = "mqttInputChannel")
public MessageHandler handler() {
return message -> {
String topic = message.getHeaders().get(MqttHeaders.RECEIVED_TOPIC).toString();
String payload = message.getPayload().toString();
log.info("收到 MQTT 消息 topic={}, payload={}", topic, payload);
try {
if (topic.endsWith("/data")) {
deviceService.handleData(payload);
} else if (topic.endsWith("/status")) {
deviceService.handleStatus(payload);
}
} catch (Exception e) {
log.error("处理失败 topic={}", topic, e);
}
};
}
/** ④ 出站通道(发送消息) */
@Bean
@ServiceActivator(inputChannel = "mqttOutboundChannel")
public MessageHandler mqttOutbound() {
MqttPahoMessageHandler handler =
new MqttPahoMessageHandler(clientId + "-out", mqttClientFactory());
handler.setAsync(true); // 异步发送
handler.setDefaultQos(1); // 默认 QoS 1
handler.setDefaultRetained(false);
return handler;
}
@Bean
public MessageChannel mqttOutboundChannel() {
return new DirectChannel();
}
}
3.3.3 发送消息(网关接口)
/** 定义网关:像调普通方法一样发 MQTT 消息 */
@MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
public interface MqttGateway {
/** 发到指定主题 */
void sendToMqtt(String data, @Header(MqttHeaders.TOPIC) String topic);
/** 发到指定主题 + 指定 QoS */
void sendToMqtt(String data,
@Header(MqttHeaders.TOPIC) String topic,
@Header(MqttHeaders.QOS) int qos);
/** 发保留消息 */
void sendRetained(String data,
@Header(MqttHeaders.TOPIC) String topic,
@Header(MqttHeaders.RETAINED) boolean retained);
}
// 使用
@Service
public class DeviceService {
@Autowired private MqttGateway mqttGateway;
public void openLock(String deviceId) {
mqttGateway.sendToMqtt("{\"cmd\":\"open\"}", "device/" + deviceId + "/cmd");
}
}
💡 另一个选择:Spring Boot 3.x / Spring Framework 6.1+ 提供了原生 MQTT 支持(
spring-mqtt),可以用@MqttListener注解,写法更像@KafkaListener。如果你的 Spring 版本够新,可以用那个,更简洁。
3.4 Python 客户端(paho-mqtt)
pip install paho-mqtt
import paho.mqtt.client as mqtt
import json
import time
# ==================== 回调(Paho 2.x API)====================
def on_connect(client, userdata, flags, reason_code, properties=None):
"""连接回调"""
if reason_code == 0:
print("连接成功")
# ★ 在 on_connect 里订阅(重连后会自动重新订阅)
client.subscribe("device/+/data", qos=1)
client.subscribe("device/+/status", qos=1)
else:
print(f"连接失败,原因码={reason_code}")
def on_message(client, userdata, msg):
"""收到消息"""
print(f"收到: {msg.topic} QoS={msg.qos} -> {msg.payload.decode()}")
try:
data = json.loads(msg.payload.decode())
handle(data)
except Exception as e:
print(f"处理失败: {e}")
def on_disconnect(client, userdata, flags, reason_code, properties=None):
print(f"断开连接,原因码={reason_code}")
def handle(data):
print(f"处理数据: {data}")
# ==================== 创建客户端 ====================
# ⚠️ Paho 2.x 必须指定 CallbackAPIVersion
client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2,
client_id="python-client-001",
clean_session=False # 保留会话,接收离线消息
)
client.username_pw_set("admin", "public")
# 设置遗嘱
client.will_set("python/client/status",
payload=json.dumps({"online": False}),
qos=1, retain=True)
client.on_connect = on_connect
client.on_message = on_message
client.on_disconnect = on_disconnect
# ==================== 连接 ====================
client.connect("localhost", 1883, keepalive=60)
client.loop_start() # ★ 启动后台网络循环线程(非阻塞)
# 或者 client.loop_forever() # 阻塞当前线程
try:
while True:
# 每 5 秒上报一次数据
payload = json.dumps({"deviceId": "D001", "temp": 26.5, "ts": int(time.time())})
client.publish("device/D001/data", payload, qos=1)
time.sleep(5)
except KeyboardInterrupt:
print("退出")
finally:
client.publish("python/client/status", json.dumps({"online": False}), qos=1, retain=True)
client.loop_stop()
client.disconnect() # 正常断开,不触发遗嘱
⚠️ Paho Python 2.x 的坑:必须传
mqtt.CallbackAPIVersion.VERSION2,否则报错。 而且回调签名变了(on_connect多了properties参数)——网上大量 1.x 的教程在 2.x 下会报错。
3.5 前端:MQTT over WebSocket
浏览器不能直接连 TCP,所以 MQTT 要走 WebSocket(端口 8083)。
npm install mqtt
import mqtt from 'mqtt'
// ★ 注意协议是 ws://(或 wss://),端口是 8083
const client = mqtt.connect('ws://localhost:8083/mqtt', {
clientId: 'web-' + Math.random().toString(16).substr(2, 8),
username: 'admin',
password: 'public',
clean: true, // 网页端一般用 true(刷新就重来)
keepalive: 30,
reconnectPeriod: 5000, // 自动重连间隔(毫秒)
connectTimeout: 10000
})
client.on('connect', () => {
console.log('MQTT 连接成功')
// 订阅
client.subscribe('device/+/status', { qos: 1 }, (err) => {
if (!err) console.log('订阅成功')
})
})
client.on('message', (topic, payload) => {
console.log(`收到 ${topic}: ${payload.toString()}`)
const data = JSON.parse(payload.toString())
// 更新 Vue / React 状态
updateDeviceStatus(topic, data)
})
client.on('reconnect', () => console.log('正在重连...'))
client.on('error', (err) => console.error('MQTT 错误', err))
// 发布(比如网页点"开灯")
function turnOnLight(deviceId) {
client.publish(`device/${deviceId}/cmd`, JSON.stringify({ cmd: 'on' }), { qos: 1 })
}
// 页面卸载时断开
window.addEventListener('beforeunload', () => client.end())
⚠️ 坑:网页端
clientId一定要随机,否则开两个标签页就会互踢。
3.6 综合案例:设备上下线 + OTA 远程升级
这是一个完整可复用的物联网场景实现。
3.6.1 主题设计
| 主题 | 方向 | 用途 | QoS | Retain |
|---|---|---|---|---|
device/{id}/status |
设备 → 云 | 上下线状态(配合遗嘱) | 1 | ✓ |
device/{id}/telemetry |
设备 → 云 | 遥测数据(温度等) | 0 | ✗ |
device/{id}/cmd |
云 → 设备 | 下发指令 | 1 | ✗ |
device/{id}/cmd/ack |
设备 → 云 | 指令执行结果回报 | 1 | ✗ |
device/{id}/ota |
云 → 设备 | OTA 固件升级通知 | 2 | ✗ |
device/{id}/ota/progress |
设备 → 云 | 升级进度 | 1 | ✗ |
3.6.2 设备端代码
/**
* 设备端:上线登记 + 接收指令 + 回报结果 + OTA 升级
*/
public class DeviceAgent {
private final String deviceId = "D001";
private MqttClient client;
public void start() throws MqttException {
client = new MqttClient("tcp://broker:1883", "device-" + deviceId);
client.setCallback(new MqttCallback() {
@Override
public void connectionLost(Throwable cause) {
log.warn("连接丢失", cause);
}
@Override
public void messageArrived(String topic, MqttMessage message) throws Exception {
String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
if (topic.equals("device/" + deviceId + "/cmd")) {
handleCommand(payload);
} else if (topic.equals("device/" + deviceId + "/ota")) {
handleOta(payload);
}
}
@Override
public void deliveryComplete(IMqttDeliveryToken token) {}
});
MqttConnectOptions opts = new MqttConnectOptions();
opts.setCleanSession(false); // 要收离线指令
opts.setAutomaticReconnect(true);
opts.setKeepAliveInterval(300); // ★ 电池设备,心跳设长一点省电
// ★ 遗嘱:意外掉线时自动标记离线
opts.setWill("device/" + deviceId + "/status",
"{\"online\":false,\"reason\":\"unexpected\"}".getBytes(), 1, true);
client.connect(opts);
// ① 订阅指令和 OTA 主题
client.subscribe("device/" + deviceId + "/cmd", 1);
client.subscribe("device/" + deviceId + "/ota", 2);
// ② ★ 主动上报"上线"(覆盖遗嘱)
publish("device/" + deviceId + "/status",
"{\"online\":true,\"fw\":\"1.0.2\"}", 1, true);
}
/** 处理云下发的指令 */
private void handleCommand(String payload) throws MqttException {
JSONObject cmd = JSON.parseObject(payload);
String cmdId = cmd.getString("cmdId"); // ★ 指令ID,用于幂等和回报
String action = cmd.getString("action");
boolean success = false;
try {
if ("open".equals(action)) {
doOpen(); // 实际开锁动作
success = true;
} else if ("close".equals(action)) {
doClose();
success = true;
}
} catch (Exception e) {
log.error("执行指令失败 cmdId={}", cmdId, e);
}
// ★ 回报执行结果(带 cmdId,云端做幂等)
String ack = JSON.toJSONString(Map.of(
"cmdId", cmdId,
"success", success,
"ts", System.currentTimeMillis()));
publish("device/" + deviceId + "/cmd/ack", ack, 1, false);
}
/** OTA 升级 */
private void handleOta(String payload) throws Exception {
JSONObject ota = JSON.parseObject(payload);
String url = ota.getString("url");
String version = ota.getString("version");
String md5 = ota.getString("md5");
// ① 下载固件
downloadFirmware(url, "/tmp/fw.bin", (percent) -> {
// ② 上报进度
publishProgress(version, percent);
});
// ③ 校验 MD5
if (!md5.equals(calcMd5("/tmp/fw.bin"))) {
publishProgress(version, -1); // -1 表示失败
return;
}
// ④ 刷写 + 重启
applyFirmware("/tmp/fw.bin");
publishProgress(version, 100);
reboot();
}
private void publish(String topic, String payload, int qos, boolean retained) throws MqttException {
MqttMessage msg = new MqttMessage(payload.getBytes(StandardCharsets.UTF_8));
msg.setQos(qos);
msg.setRetained(retained);
client.publish(topic, msg);
}
/** 正常下线:先发离线状态,再断开 */
public void shutdown() throws MqttException {
publish("device/" + deviceId + "/status",
"{\"online\":false,\"reason\":\"normal\"}", 1, true);
client.disconnect();
}
}
3.6.3 云端代码(下发指令 + 等回报)
/**
* 云端:下发指令并等待设备回报(★ 应用层可靠性设计)
*/
@Service
public class CommandService {
@Autowired private MqttClientService mqtt;
@Autowired private StringRedisTemplate redis;
/**
* 下发指令,同步等待设备回报(最多 waitSeconds 秒)
*/
public CommandResult sendCommand(String deviceId, String action, int waitSeconds) {
String cmdId = UUID.randomUUID().toString();
// ① 在 Redis 里登记"等待回报"(key = cmdId)
String waitKey = "mqtt:cmd:" + cmdId;
redis.opsForValue().set(waitKey, "WAITING", Duration.ofSeconds(waitSeconds + 5));
// ② 下发指令(QoS 1 保证送达)
String payload = JSON.toJSONString(Map.of("cmdId", cmdId, "action", action));
mqtt.publish("device/" + deviceId + "/cmd", payload, 1, false);
// ③ 轮询等回报(生产上更推荐:MQ 消费者异步处理 + WebSocket 推前端)
long deadline = System.currentTimeMillis() + waitSeconds * 1000L;
while (System.currentTimeMillis() < deadline) {
String result = redis.opsForValue().get(waitKey);
if (!"WAITING".equals(result)) {
return JSON.parseObject(result, CommandResult.class);
}
Thread.sleep(200);
}
return CommandResult.timeout(cmdId);
}
/**
* 订阅 cmd/ack 主题,处理设备回报
*/
@MqttListener(topics = "device/+/cmd/ack", qos = "1") // Spring MQTT 6.x 写法
public void onAck(String payload) {
JSONObject ack = JSON.parseObject(payload);
String cmdId = ack.getString("cmdId");
String waitKey = "mqtt:cmd:" + cmdId;
// ★ 幂等:只有还在 WAITING 才处理(防止设备重复回报)
if ("WAITING".equals(redis.opsForValue().get(waitKey))) {
redis.opsForValue().set(waitKey, payload, Duration.ofMinutes(5));
} else {
log.warn("收到重复或过期的回报 cmdId={}", cmdId);
}
}
}
★ 这个案例体现的三个关键设计:
- 遗嘱 + 主动上报组合,准确感知设备上下线。
- 指令 ID(cmdId) 做幂等,杜绝重复执行。
- 应用层回报机制弥补 QoS 的不足——QoS 只保证“消息送达”,不保证“设备执行成功”。
第四章:应用场景、生产问题与解决方案、安全与性能
这一章是实战价值最高的部分。4.2 节列出 15 个线上真实会遇到的问题,每个都给「现象 → 原因 → 解决方案」。
4.1 五个典型应用场景
场景一:车联网(最典型)
业务:车辆上报 GPS、里程、油量、故障码;云端下发远程控制(解锁、鸣笛、限功率)。
车载 T-Box ──MQTT(QoS1)──► EMQX 集群 ──► Kafka ──► 实时计算(Flink)──► 监控大屏
▲ │ │
└────── 远程指令(QoS1) ───┘ │
└──► 时序数据库(TDengine/InfluxDB)───────────┘
| 需求 | MQTT 方案 |
|---|---|
| 高频 GPS 上报 | QoS 0 或 1,每秒 1 条 |
| 远程控制指令 | QoS 1 + 应用层回报(cmdId) |
| 车辆离线检测 | 遗嘱 + retained 状态 |
| 百万车辆连接 | EMQX 集群 + 负载均衡 |
| 数据要分析 | 规则引擎转 Kafka |
场景二:智能家居
业务:手机 App 控制灯、空调、窗帘;设备上报状态。
手机 App ──MQTT(ws)──► Broker ◄──MQTT── 智能灯 / 空调 / 传感器
│
└──► 音箱 / 云端场景联动
| 需求 | MQTT 方案 |
|---|---|
| App 秒开看到当前状态 | ★ 保留消息(设备状态 retain=true) |
| 多手机同时控制 | 广播特性天然支持 |
| 设备离线提醒 | 遗嘱消息 |
| 局域网也能用 | 本地 Broker + 云端桥接 |
场景三:工业物联网(IIoT)
业务:工厂设备数据采集(PLC、传感器)、设备预测性维护、产线监控。
| 特点 | 应对 |
|---|---|
| 数据点极多(单厂上万点) | 主题分层设计 factory/{厂}/line/{线}/device/{设备}/{指标} |
| 网络不稳定(车间干扰) | QoS 1 + 断线重连 + 离线缓冲 |
| 实时性要求高 | 边缘网关本地处理 + MQTT 上云 |
| 历史数据要存 | MQTT → Kafka → 时序库 |
场景四:共享单车 / 共享充电宝
业务:设备上报位置和电量;云端下发开锁指令。
| 需求 | 方案 |
|---|---|
| 省电(设备靠电池) | ★ KeepAlive 设长(300~1800 秒),减少心跳 |
| 开锁要可靠 | QoS 1 + 应用层回报 + 超时重试 |
| 海量设备(百万级) | EMQX 集群,ClientId = 设备序列号 |
| 离线也要能开(蓝牙) | MQTT 为主,蓝牙为兜底 |
场景五:移动 App 推送 / IM
业务:App 消息推送、聊天室、直播弹幕。
| 需求 | 方案 |
|---|---|
| App 后台也要收 | MQTT 长连接(比各家推送 SDK 更可控) |
| 离线消息 | CleanSession=false + QoS 1 |
| 群聊/直播间 | 订阅同一主题,广播特性天然适合 |
| 海量在线 | EMQX 支持百万连接 |
⚠️ 现实提醒:移动端 App 推送在国内通常用厂商推送通道(小米/华为/OPPO 推送)+ 个推/极光,因为 App 被杀后台后长连接也断了。MQTT 更适合App 在前台时的实时消息。
4.2 ★ 15 个生产常见问题与解决方案
问题 1:消息丢失
现象:设备发了数据,云端没收到。
原因分析(按概率排序):
| 原因 | 说明 |
|---|---|
| ① 用了 QoS 0 | 最常见。QoS 0 就是“发完不管”,网络一抖就丢 |
| ② 订阅端 QoS 比发布端低 | QoS 取小值,发布 QoS1 + 订阅 QoS0 = 实际 QoS0 |
| ③ 设备离线且 CleanSession=true | 离线期间的消息直接丢弃 |
| ④ 离线队列满了 | Broker 有队列上限,超了丢弃最老的 |
| ⑤ Broker 重启/崩溃 | 非持久化的消息丢失 |
解决方案:
// ① 关键数据用 QoS 1(至少一次)
client.publish("device/D001/alarm", payload, 1, false);
// ② 发布和订阅两端都要设 QoS(★ 取小值陷阱)
client.subscribe("device/+/alarm", 1); // 订阅端也要 QoS 1
// ③ 需要离线消息:CleanSession=false + ClientId 固定
opts.setCleanSession(false);
opts.setClientId("device-" + deviceSn); // ★ 必须固定,否则找不到会话
// ④ 应用层兜底:设备本地缓存 + 重发
// 设备端把每条消息先写本地队列,收到 PUBACK 才删除
问题 2:消息重复(QoS 1 的经典代价)
现象:同一条数据被处理了两次,比如订单重复扣款、告警重复发。
原因:QoS 1 是“至少一次“。发送方没收到 PUBACK 就会重发,Broker 可能已经收到了 → 投递两次。
解决方案:★ 消费端必须做幂等(这是铁律,不是可选项)。
/**
* 方案一:消息 ID 去重(推荐)
*/
@Service
public class MqttMessageHandler {
@Autowired private StringRedisTemplate redis;
public void handle(String topic, String payload) {
JSONObject data = JSON.parseObject(payload);
// ★ 设备端上报时带上唯一 msgId
String msgId = data.getString("msgId");
if (msgId == null) {
log.warn("消息缺少 msgId,无法去重, topic={}", topic);
return;
}
// SETNX:能设置成功 = 第一次处理;已存在 = 重复消息,丢弃
Boolean first = redis.opsForValue()
.setIfAbsent("mqtt:dedup:" + msgId, "1", Duration.ofHours(24));
if (Boolean.FALSE.equals(first)) {
log.debug("重复消息,丢弃 msgId={}", msgId);
return;
}
// 真正的业务处理
doBusiness(data);
}
}
/**
* 方案二:数据库唯一索引兜底(最可靠)
*/
// 建表时给业务唯一键加唯一索引
// CREATE UNIQUE INDEX uk_msg_id ON device_data(msg_id);
// 插入时重复会抛 DuplicateKeyException,catch 掉即可
try {
deviceDataMapper.insert(data);
} catch (DuplicateKeyException e) {
log.debug("重复数据,忽略 msgId={}", data.getMsgId());
}
/**
* 方案三:业务状态机(最本质的幂等)
*/
// 比如"关单"操作:先查状态,已经是"已关闭"就直接返回成功,不重复执行
if ("CLOSED".equals(order.getStatus())) {
return Result.ok("已关闭"); // 幂等
}
★ 面试话术:“QoS 1 保证的是至少一次,重复是协议允许的。所以消费端幂等是必须的,我们用的是’设备端生成 msgId + Redis SETNX 去重’ + ’数据库唯一索引兜底’两层防护。”
问题 3:消息乱序
现象:设备先发 A 后发 B,云端先收到 B 后收到 A。
原因:
- QoS 0 没有顺序保证(网络重传、多路径)。
- QoS 1/2 只保证单条消息的顺序,但如果上一条还在“飞行中”(等待 PUBACK),下一条可能先到。
- 多个连接并发发布。
解决方案:
// ① 设备端加序列号,云端按序处理
// 设备发送时带 seq(递增),云端缓存并按 seq 排序后处理
public void handle(String payload) {
JSONObject data = JSON.parseObject(payload);
long seq = data.getLongValue("seq");
long expected = seqHolder.get(deviceId) + 1;
if (seq < expected) {
return; // 旧消息,丢弃
} else if (seq > expected) {
buffer.put(seq, data); // 未来的消息,先缓存
return;
}
process(data); // 正好是下一条,处理
seqHolder.set(deviceId, seq);
// 再把 buffer 里连续的补上
flushBuffer(deviceId);
}
// ② 关键场景用 QoS 2(严格保序)—— 但性能差
// ③ 最简单:业务上不要求严格顺序(如温度上报,乱序无所谓)
⚠️ 诚实提醒:MQTT 不保证跨消息的全局顺序。如果你的业务强依赖顺序(如交易流水),要么在应用层用 seq 重排,要么别用 MQTT。
问题 4:设备明明断了,但状态还是“在线”(假在线)
现象:设备断电了,监控大屏还显示在线。
原因:
- 没设遗嘱,或遗嘱主题/内容写错。
- KeepAlive 设太长(如 1800 秒),Broker 要 45 分钟才判定离线。
- Broker 没检测到 TCP 断开(网络中间设备保持连接)。
解决方案:
// ① 必设遗嘱
opts.setWill("device/" + id + "/status",
"{\"online\":false}".getBytes(), 1, true);
// ② KeepAlive 按业务容忍度设(实时监控就设短)
opts.setKeepAliveInterval(60); // 实时监控:60 秒
// opts.setKeepAliveInterval(300); // 电池设备:300 秒(代价是掉线发现慢)
// ③ ★ 服务端主动探活兜底:定时查"最后上报时间"
@Scheduled(fixedRate = 60000)
public void checkStaleDevices() {
// 查最后心跳超过 3 分钟的设备,标记为"疑似离线"并告警
List<Device> stale = deviceMapper.selectByLastSeenBefore(
LocalDateTime.now().minusMinutes(3));
stale.forEach(d -> {
deviceMapper.updateOnline(d.getId(), false);
alertService.send("设备 " + d.getId() + " 长时间无数据上报");
});
}
⚠️ 坑:正常
disconnect()不触发遗嘱!所以设备正常关机时要自己发一条 offline(retained=true),否则状态永远卡在 online。见 3.6.2 的shutdown()。
问题 5:ClientId 重复导致“互踢”
现象:设备 A 和 B 反复掉线重连,日志里全是“Connection lost”。
原因:两个(或多个)连接用了同一个 ClientId,Broker 只允许一个,后连的踢掉先连的,先连的重连又踢后连的 → 死循环。
解决方案:
// ❌ 错误(抄 Demo 抄出来的)
opts.setClientId("mqtt_client");
// ✅ 正确:用设备唯一标识
opts.setClientId("device-" + deviceSn); // 设备序列号
// 或
opts.setClientId("device-" + mac.replace(":", ""));
// ⚠️ 如果确实需要随机 ClientId(如网页端),记得 CleanSession=true
opts.setClientId("web-" + UUID.randomUUID());
opts.setCleanSession(true); // 随机的 ClientId 不需要保留会话
排查方法:
# EMQX 控制台 → 连接管理,看是否有同 ClientId 反复上下线
# 或用 CLI
emqx ctl clients list | grep D001
问题 6:遗嘱不触发 / 误触发
| 现象 | 原因 | 解决 |
|---|---|---|
| 设备断电了,遗嘱没发 | ① 没设 setWill;② Broker 也挂了;③ KeepAlive 太长还没超时 |
检查 CONNECT 参数;缩短 KeepAlive |
| 设备正常下线,遗嘱却发了 | 用了 client.close() 而不是 disconnect(),或直接进程 kill 导致 TCP 异常断 |
正常路径必须调 disconnect() |
| 遗嘱发了但没人收到 | 订阅主题写错 / 订阅端 QoS 0 降级 | 检查主题拼写和订阅 QoS |
// ✅ 正确的正常下线流程
public void shutdown() {
// ① 先主动发离线状态(覆盖 retained)
client.publish("device/D001/status",
"{\"online\":false}".getBytes(), 1, true);
// ② 再正常断开(不会触发遗嘱)
client.disconnect();
}
问题 7:保留消息变成“脏数据”
现象:设备早就下线了,但新订阅者一上来就收到一个很老的状态(retained 消息没清理)。
原因:保留消息永久保存在 Broker(按主题只留最新一条),设备异常下线没来得及清理。
解决方案:
// ① 设备下线时主动清空保留消息:发空载荷 + retained=true
client.publish("device/D001/status", new byte[0], 1, true);
// ② MQTT 5.0:给保留消息设过期时间(Message Expiry Interval)
// 过期后 Broker 自动删除
// ③ 服务端定期巡检:超过 N 天没更新的 retained 消息清掉
@Scheduled(cron = "0 0 4 * * ?")
public void cleanStaleRetained() {
// 用 EMQX HTTP API 列出 retained 消息,检查时间戳,清理老的
}
问题 8:Broker 内存暴涨 / 僵尸会话堆积
现象:EMQX 运行几个月后内存持续增长,重启才恢复。
原因:
- 僵尸会话:设备刷机/报废后再也不上线,但
CleanSession=false让会话永久保留。 - 离线队列积压:设备长期离线,给它暂存的消息越堆越多。
解决方案:
# ① MQTT 5.0:设会话过期时间(★ 最根本的解法)
session-expiry-interval: 86400 # 会话保留 1 天,过期自动清理
# ② EMQX 配置:限制离线队列长度
# emqx.conf
mqtt.max_queued_messages = 10000 # 每个会话最多暂存 1 万条
mqtt.max_inflight = 100 # 最多 100 条在飞未确认
# ③ 3.1.1 场景:定期清理长期不活跃的会话(EMQX HTTP API)
curl -X DELETE http://localhost:18083/api/v5/clients/device-old-001
// ④ 服务端兜底:定时清理超过 7 天没上报的设备会话
@Scheduled(cron = "0 0 3 * * ?")
public void cleanZombieSessions() {
List<Device> dead = deviceMapper.selectOfflineSince(LocalDateTime.now().minusDays(7));
dead.forEach(d -> emqxApi.kickAndCleanSession("device-" + d.getSn()));
}
问题 9:回调线程阻塞导致心跳超时,连接频繁断开
现象:连接每隔几分钟就断开重连一次,日志显示 Connection lost。
原因:★ Paho 的回调和心跳在同一个网络线程里。你在 messageArrived() 里做了耗时操作(调慢接口、同步等锁、写数据库慢),心跳发不出去 → Broker 判定超时 → 断开。
// ❌ 错误:在回调里做耗时操作
@Override
public void messageArrived(String topic, MqttMessage message) {
String payload = new String(message.getPayload());
// 同步调用一个 3 秒的外部接口 + 慢 SQL
externalService.callSlow(payload); // 阻塞 3 秒!
slowMapper.insert(parse(payload)); // 又阻塞 1 秒
// → 心跳发不出去 → 被 Broker 踢掉
}
// ✅ 正确:回调里只做投递,业务丢到线程池异步处理
@Override
public void messageArrived(String topic, MqttMessage message) {
String payload = new String(message.getPayload());
// ★ 立刻把消息丢进队列/线程池,毫秒级返回
businessExecutor.submit(() -> {
try {
externalService.callSlow(payload);
slowMapper.insert(parse(payload));
} catch (Exception e) {
log.error("处理失败", e);
}
});
}
// 线程池配置
@Bean("businessExecutor")
public ExecutorService businessExecutor() {
return new ThreadPoolExecutor(
8, 32, 60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(10000), // ★ 有界队列,防止 OOM
new ThreadFactoryBuilder().setNameFormat("mqtt-biz-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy());
}
⚠️ 另外注意:回调里必须 try-catch,异常抛出去会导致 Paho 线程崩溃、连接断开。
问题 10:网络闪断后重连风暴
现象:Broker 重启或网络抖动后,10 万设备同时重连,把 Broker 冲垮(连接风暴)。
原因:所有设备的重连时间都在同一瞬间。
解决方案:★ 重连加随机退避(Jitter)。
/**
* 自定义重连策略:指数退避 + 随机抖动
*/
public class ReconnectWithJitter implements Runnable {
private final MqttClient client;
private int attempt = 0;
private static final int MAX_BACKOFF_SEC = 300; // 最大退避 5 分钟
@Override
public void run() {
while (!client.isConnected()) {
try {
// ① 指数退避:2^n 秒
long backoff = Math.min((long) Math.pow(2, attempt), MAX_BACKOFF_SEC);
// ② ★ 加随机抖动:±30%,打散重连时间
long jitter = (long) (backoff * (0.7 + Math.random() * 0.6));
log.info("第 {} 次重连,等待 {} 秒", attempt + 1, jitter);
Thread.sleep(jitter * 1000);
client.reconnect();
attempt = 0; // 成功则重置计数
log.info("重连成功");
} catch (Exception e) {
attempt++;
log.warn("重连失败,attempt={}", attempt, e);
}
}
}
}
# EMQX 侧也能限速,防止被打爆
# emqx.conf
listeners.tcp.default.max_connections = 1000000
listeners.tcp.default.max_conn_rate = 1000 # ★ 每秒最多接受 1000 个新连接
问题 11:消息体太大被拒
现象:发一张图片或一个大 JSON,Broker 拒绝或断开连接。
原因:MQTT 协议最大报文 256MB,但 Broker 通常会配置更小的限制(EMQX 默认 1MB)。
解决方案:
# ① 调大限制(不推荐,大报文会阻塞连接)
# emqx.conf
mqtt.max_packet_size = 10MB
// ② ★ 正确做法:大文件不要走 MQTT,走"URL + 对象存储"
// 步骤:设备上传文件到 OSS → 得到 URL → MQTT 只发 URL
String ossUrl = ossClient.upload(file); // 文件走 HTTP 传到对象存储
client.publish("device/D001/photo", // MQTT 只发一个几十字节的 URL
ossUrl.getBytes(), 1, false);
// ③ 分片:把大 payload 拆成多条小消息,带 seq 让接收方组装
// payload: {"seq":1,"total":10,"data":"..."}
★ 原则:MQTT 适合小报文高频。超过几十 KB 的消息都应该重新设计。
问题 12:主题设计混乱,后期无法维护
现象:主题五花八门(Device001/temp、device/001/Temp、dev/001/t),ACL 没法配,通配符订阅失效。
解决方案:★ 制定主题规范并写进文档。
推荐格式:{产品}/{设备ID}/{数据类型}
示例:
vehicle/VIN123456/gps 车辆 GPS
vehicle/VIN123456/status 车辆状态
device/D001/telemetry 设备遥测
device/D001/cmd 云端指令(下行)
device/D001/cmd/ack 指令回报(上行)
规范:
① 全小写,不用空格
② 不以 / 开头
③ 设备 ID 放固定层级,便于通配符和 ACL
④ 上行(设备→云)和下行(云→设备)用不同后缀区分(cmd / cmd/ack)
⑤ 层级不超过 5 层
ACL 配合示例:
-- EMQX 内置数据库认证:只允许设备发布自己的主题,订阅自己的指令
-- 设备 D001 只能:
-- publish : device/D001/# (除了 cmd)
-- subscribe: device/D001/cmd
问题 13:安全——Broker 被匿名访问 / 数据被窃听
现象:Broker 暴露在公网,任何人都能连上来订阅所有主题,甚至发布控制指令。
原因:默认配置允许匿名连接,且用明文 TCP(1883)。
解决方案(★ 生产必做的五件事):
# ① 关闭匿名认证(EMQX)
# emqx.conf
allow_anonymous = false
# ② 每个设备独立账号(不要所有设备共用一个账号)
-- ③ 建认证表(EMQX MySQL 认证)
CREATE TABLE mqtt_user (
username VARCHAR(100) PRIMARY KEY,
password_hash VARCHAR(255) NOT NULL, -- ★ 存 bcrypt 哈希,绝不存明文
salt VARCHAR(64),
is_superuser TINYINT DEFAULT 0,
created_at DATETIME
);
-- ④ 建 ACL 表:限制每个设备只能操作自己的主题
CREATE TABLE mqtt_acl (
username VARCHAR(100) NOT NULL,
permission VARCHAR(10) NOT NULL, -- allow / deny
action VARCHAR(10) NOT NULL, -- publish / subscribe / all
topic VARCHAR(255) NOT NULL -- 支持 %u 占位符(用户名)、%c(ClientId)
);
-- 例:设备 D001(username=device_D001)只能发布自己的数据、订阅自己的指令
INSERT INTO mqtt_acl VALUES
('device_D001', 'allow', 'publish', 'device/D001/telemetry'),
('device_D001', 'allow', 'publish', 'device/D001/status'),
('device_D001', 'allow', 'subscribe', 'device/D001/cmd'),
('device_D001', 'deny', 'all', 'device/#'); -- ★ 拒绝访问其他设备
// ⑤ 强制 TLS(客户端侧)
MqttClient client = new MqttClient("ssl://broker.example.com:8883", clientId);
// 配合 3.2.4 的 SSLContext 配置
| 安全措施 | 做法 |
|---|---|
| 关闭匿名 | allow_anonymous = false |
| 一机一密 | 每个设备独立账号,密码 bcrypt 哈希存储 |
| TLS 加密 | 用 8883 端口,配 CA 证书 |
| ACL 权限 | 限制每个账号只能操作自己的主题 |
禁用 # 订阅 |
EMQX 可配规则,禁止客户端订阅 # |
| 速率限制 | 限制每个客户端的发布速率,防止刷爆 |
问题 14:数据要长期存储和分析,MQTT 存不住
现象:Broker 里的消息消费完就没了,没法做历史查询和大数据分析。
原因:★ MQTT 不是存储系统,它是“传完就扔”的管道。
解决方案:★ 规则引擎转 Kafka / 数据库(EMQX 内置能力)。
-- EMQX 规则引擎 SQL:把设备数据转发到 Kafka
SELECT
payload.temp AS temp,
payload.humidity AS humidity,
topic,
timestamp AS ts
FROM "device/+/telemetry"
WHERE payload.temp > 100 -- 还能过滤
数据流:
设备 ──MQTT──► EMQX ──规则引擎──► Kafka ──► Flink 实时计算 ──► 告警
│ └──► 数仓(离线分析)
└──► TDengine / InfluxDB(时序数据,直接存)
└──► MySQL / Redis(设备最新状态)
★ 架构原则:MQTT 负责“连得上”,Kafka 负责“存得住、算得动”。不要指望 MQTT 做持久化。
问题 15:跨地域 / 多机房如何同步
现象:设备分布在全国,连到不同区域的 Broker,云端要统一处理。
解决方案:MQTT 桥接(Bridge)。
# EMQX 桥接配置:把本地 Broker 的消息转发到中心 Broker
bridges:
mqtt:
central:
enable: true
server: "broker-central.example.com:1883"
username: bridge_user
password: "xxx"
# 把本地的 device/+/telemetry 转发到中心,加前缀
forwards:
- topic: "device/+/telemetry"
remote_topic: "region-east/${topic}"
qos: 1
# 从中心订阅下行指令
subscriptions:
- topic: "cmd/region-east/#"
qos: 1
边缘 Broker(华东)──桥接──► 中心 Broker(北京)──► Kafka
边缘 Broker(华南)──桥接──► │
└──► 统一下发指令
⚠️ 坑:桥接是异步转发,不是强一致。桥接断开期间的消息可能丢失(取决于 Broker 是否做桥接侧的持久化)。关键指令要有应用层确认。
4.3 MQTT 安全清单(★ 生产必做)
| # | 措施 | 优先级 | 说明 |
|---|---|---|---|
| 1 | 关闭匿名访问 | P0 | allow_anonymous = false |
| 2 | 一机一密 | P0 | 每个设备独立账号,禁止共用 |
| 3 | 密码哈希存储 | P0 | bcrypt / PBKDF2,绝不存明文 |
| 4 | 启用 TLS | P0 | 端口 8883,配 CA 证书;内网可豁免 |
| 5 | 配置 ACL | P0 | 设备只能访问自己的主题 |
| 6 | 禁止 # 订阅 |
P1 | 防止一个客户端收走所有数据 |
| 7 | 速率限制 | P1 | 防止恶意设备刷消息 |
| 8 | 定期轮换凭据 | P1 | 设备密钥定期更新(OTA 下发新密钥) |
| 9 | 审计日志 | P2 | 记录连接、订阅、发布行为 |
| 10 | 网络隔离 | P2 | Broker 不直接暴露公网,走网关/VPN |
📖 更全面的物联网与云安全内容见
13-新兴攻击面与专项安全/第五章(实时通信与 IoT 协议安全)。
4.4 性能优化与海量连接
单机能撑多少连接?
| Broker | 单机连接数(参考) | 说明 |
|---|---|---|
| EMQX | 百万级 | Erlang/OTP,天然高并发 |
| Mosquitto | 万级 | C 写的,轻量但扩展性弱 |
| HiveMQ | 十万~百万级 | Java,商业版强 |
调优清单
| 优化项 | 做法 | 效果 |
|---|---|---|
| 用 EMQX 而非 Mosquitto | 生产选 EMQX | ★★★★★ |
| 心跳设长 | 电池设备 300~1800 秒 | ★★★★★ 省电省带宽 |
| QoS 降级 | 非关键数据用 QoS 0 | ★★★★☆ 减少确认报文 |
避免 # 订阅 |
用精确主题 | ★★★★☆ 减少 Broker 匹配开销 |
| 主题别名(5.0) | 用数字代替长主题名 | ★★★☆☆ 省流量 |
| 批量上报 | 设备攒 10 条发一次 | ★★★★☆ 减少报文数 |
| Payload 压缩/精简 | JSON 改二进制(Protobuf/MessagePack) | ★★★★☆ |
| 会话过期 | 设 session-expiry,清理僵尸会话 | ★★★★☆ |
| 限制队列长度 | max_queued_messages |
★★★☆☆ 防内存爆 |
| 集群 + 负载均衡 | EMQX 集群 + LB | ★★★★★ |
| 连接速率限制 | max_conn_rate |
★★★☆☆ 防重连风暴 |
系统参数调优(Linux)
# ① 文件描述符(每个连接占一个 fd)
ulimit -n 1000000
echo "* soft nofile 1000000" >> /etc/security/limits.conf
echo "* hard nofile 1000000" >> /etc/security/limits.conf
# ② TCP 参数
sysctl -w net.core.somaxconn=32768
sysctl -w net.ipv4.tcp_max_syn_backlog=32768
sysctl -w net.ipv4.ip_local_port_range="1024 65535"
sysctl -w net.core.netdev_max_backlog=32768
# ③ EMQX 的 Erlang 进程数
# emqx.conf
node.max_ports = 1000000
Payload 精简:JSON → 二进制
// ❌ JSON:{"deviceId":"D001","temp":26.5,"hum":60,"ts":1700000000} ≈ 55 字节
// ✅ Protobuf / MessagePack:≈ 12 字节(省 78%)
// 或者最简单:用数组代替 key
// ["D001",26.5,60,1700000000] ≈ 30 字节
// ★ 对百万设备每秒上报的场景,这 25 字节的差距 = 每天省几十 GB 流量
第五章:选型对比、面试题与小结
5.1 MQTT vs 其他协议:怎么选
| 协议 | 传输层 | 模式 | 报文开销 | 双向推送 | 堆积能力 | 典型场景 |
|---|---|---|---|---|---|---|
| MQTT | TCP | 发布/订阅 | ★ 极小(2 字节起) | ✅ | 弱 | 物联网、移动推送 |
| HTTP | TCP | 请求/响应 | 大(几百字节) | ❌ | 无 | Web API |
| HTTP/2 + SSE | TCP | 请求/响应 + 服务端推 | 中 | ✅(SSE) | 无 | Web 实时推送 |
| WebSocket | TCP | 全双工 | 中 | ✅ | 无 | 浏览器实时通信 |
| CoAP | UDP | 请求/响应 | ★ 极小 | ✅ | 无 | 极受限设备(传感器) |
| AMQP | TCP | 队列/订阅 | 中 | ✅ | 强 | 企业级消息(RabbitMQ) |
| Kafka 协议 | TCP | 发布/订阅(分区) | 中 | ✅ | ★ 极强 | 大数据流、日志 |
和 CoAP 的区别(面试偶尔问)
| MQTT | CoAP | |
|---|---|---|
| 传输层 | TCP(可靠、有连接) | UDP(不可靠、无连接) |
| 模式 | 发布/订阅(长连接) | 请求/响应(像 HTTP) |
| 开销 | 小 | ★ 更小 |
| 双向推送 | ✅ 天然支持 | 需 observe 机制 |
| NAT 穿透 | 需长连接保持 | UDP 打洞更麻烦 |
| 适用 | 大部分 IoT 场景 | 电池供电的极简传感器 |
结论:绝大多数物联网场景用 MQTT。CoAP 只在“设备资源极端受限 + 只需偶尔上报一次”的场景(如智能水表)才考虑。
★ 和 Kafka/RocketMQ 的配合关系(不是替代!)
【错误认知】"用 MQTT 还是 Kafka?" —— 这是伪命题,它们是上下游关系。
【正确架构】
设备 ──MQTT──► Broker ──规则引擎──► Kafka ──► 业务系统
(海量长连接) (协议转换) (堆积/削峰/持久化)
MQTT:解决"百万设备连得上"
Kafka:解决"海量数据存得住、算得动"
5.2 Broker 选型
| Broker | 语言 | 单机连接 | 集群 | 特点 | 适用 |
|---|---|---|---|---|---|
| EMQX | Erlang | ★ 百万级 | ✅ 原生集群 | 国产、中文文档全、规则引擎强大、有企业版 | ★ 国内首选 |
| Mosquitto | C | 万级 | ❌(需桥接) | 极轻量、简单 | 学习、小项目、边缘设备 |
| HiveMQ | Java | 十万~百万 | ✅ | 商业版强、MQTT 5 支持好 | 欧美企业 |
| VerneMQ | Erlang | 十万级 | ✅ | 开源、可水平扩展 | 中等规模 |
| 云厂商 IoT 平台 | — | 弹性 | ✅ | 阿里云/华为云/AWS IoT | ★ 不想自己运维就选这个 |
选型建议:
| 你的情况 | 选什么 |
|---|---|
| 国内项目、要自己运维 | EMQX(社区版免费,功能足够) |
| 学习 / Demo / 边缘网关 | Mosquitto |
| 设备量 < 1 万、不想运维 Broker | 阿里云 IoT / 华为云 IoT 等托管服务 |
| 设备量 > 100 万、有多地域 | EMQX 企业版 或 云厂商 IoT 平台 |
💡 成本提示:自己运维 EMQX 要管集群、监控、升级、扩容。如果设备量不大(<10 万),云厂商 IoT 平台按连接数计费可能更划算,还省了运维人力。
5.3 面试题(共 25 题)
5.3.1 基础题(1~10,★★☆☆☆)
Q1. 什么是 MQTT?和 HTTP 的区别? ★★☆☆☆
MQTT 是基于 TCP 的发布/订阅模式的超轻量物联网协议。 和 HTTP 的区别:① MQTT 是长连接,HTTP 是请求/响应;② MQTT 服务端能主动推,HTTP 不能;③ MQTT 头部极小(2 字节起),HTTP 几百字节;④ MQTT 为弱网、低功耗、海量设备设计。
Q2. MQTT 的三个 QoS 级别? ★★★☆☆
- QoS 0 最多一次:发完不管,可能丢。最快最省。
- QoS 1 至少一次:PUBACK 确认,保证送达但可能重复。
- QoS 2 恰好一次:PUBREC→PUBREL→PUBCOMP 四次握手,不丢不重,开销最大。 ★ 面试加分:说出“发布 QoS 和订阅 QoS 取小值”这个降级规则。
Q3. QoS 1 消息会重复吗?怎么办? ★★★★☆
会。QoS 1 只保证“至少一次”。没收到 PUBACK 就重发,可能导致重复投递。 必须在消费端做幂等:
msgId+ Redis SETNX 去重,或数据库唯一索引兜底,或业务状态机。
Q4. 什么是保留消息? ★★☆☆☆
Broker 记住每个主题的最后一条 retained 消息,新订阅者一上线立刻收到。用于“设备当前状态”(灯是开是关)。 清除方法:发一条空载荷 + retained=true。
Q5. 什么是遗嘱消息(LWT)?什么时候触发? ★★★☆☆
设备 CONNECT 时预先登记的“遗言”,它异常掉线时 Broker 自动代发。用于设备离线检测。 ★ 正常调
disconnect()不触发遗嘱——只有 TCP 异常断或心跳超时才触发。
Q6. 发布/订阅模式的解耦体现在哪? ★★☆☆☆
发布者不知道谁会收到,订阅者不知道谁发的,双方通过 Topic 和 Broker 间接通信。空间解耦(不知道对方地址)、时间解耦(不需要同时在线)。
Q7. CleanSession 是什么意思? ★★★☆☆
true= Broker 不保留会话,断开即忘,离线消息丢失,重连要重新订阅。false= Broker 保留会话(订阅关系 + 未确认消息),重连后自动恢复订阅并补发离线消息(仅 QoS 1/2)。 ★ 要收离线消息必须CleanSession=false+ ClientId 固定。
Q8. MQTT 的心跳机制? ★★★☆☆
客户端 CONNECT 时声明 KeepAlive 秒数,必须在 1.5×KeepAlive 内发至少一个报文(无数据就发 PINGREQ),Broker 回 PINGRESP。超时未收到 → 判定掉线 → 触发遗嘱。 ★ 权衡:心跳短=掉线发现快但耗电;电池设备建议 300~1800 秒。
Q9. MQTT 主题通配符? ★★☆☆☆
+匹配单层(home/+/temp),#匹配多层且必须在最后(home/#)。 ⚠️ 只有订阅时能用,发布时用了会发布到字面上带通配符的主题。 ⚠️#不会匹配$SYS/...系统主题。
Q10. MQTT 默认端口? ★☆☆☆☆
1883(TCP 明文)、8883(TLS)、8083(WebSocket)、8084(WSS)、18083(EMQX 控制台)。
5.3.2 进阶题(11~18,★★★☆☆ ~ ★★★★☆)
Q11. QoS 2 的四次握手过程?为什么能保证恰好一次? ★★★★☆
PUBLISH→PUBREC→PUBREL→PUBCOMP。 发送方收到 PUBREC 后记住消息已送达但未完成;Broker 在收到 PUBREL 前不投递给订阅者,收到后只投递一次。双方都用 msgId 去重,所以重传也不会重复投递。 加分点:“即使 QoS 2 也只保证消息层面的恰好一次,业务层面还要靠指令 ID + 执行回报。”
Q12. 设备频繁掉线重连,怎么排查? ★★★★☆
① ClientId 重复互踢(最常见)——检查是否用了固定字符串当 ClientId。 ② 回调线程阻塞导致心跳发不出——检查
messageArrived()里有没有耗时操作。 ③ KeepAlive 太短而网络延迟大。 ④ Broker 连接数打满。 ⑤ 认证失败(CONNACK 返回码 4/5)。 加分点:“重连要加指数退避 + 随机抖动,否则设备多了会形成重连风暴把 Broker 冲垮。”
Q13. 设备断电了,为什么状态还是在线? ★★★★☆
① 没设遗嘱;② KeepAlive 太长,Broker 还没判定超时;③ Broker 没检测到 TCP 断开。 解决:设遗嘱 + 合理 KeepAlive + 服务端兜底巡检(定时查最后上报时间,超阈值判离线)。 加分点:“正常下线要主动发 offline 再
disconnect(),因为正常断开不触发遗嘱。”
Q14. MQTT 能保证消息顺序吗? ★★★★☆
★ 不能完全保证。QoS 0 无序;QoS 1/2 在同一连接、同一 msgId 的生命周期内有顺序,但前一条在“飞行中”时后一条可能先到;多连接并发更无法保证。 解决:应用层加 seq 序列号重排;或关键场景用 QoS 2;或业务上不依赖严格顺序。
Q15. 百万设备连接,Broker 怎么扛? ★★★★☆
① 选 EMQX(Erlang 高并发,单机百万连接);② 集群 + 负载均衡;③ 调大 Linux
ulimit -n;④ 心跳设长减少报文;⑤ QoS 降级;⑥ 避免#订阅;⑦ 限速防重连风暴;⑧ 会话过期清理僵尸会话。 加分点:“MQTT 只负责连得上,数据要转 Kafka 才能存得住算得动。”
Q16. MQTT 怎么做设备权限控制? ★★★☆☆
① 关闭匿名访问;② 一机一密(每设备独立账号,密码 bcrypt 哈希);③ ACL 限制每个账号只能发布/订阅自己的主题;④ 启用 TLS;⑤ 禁止
#订阅;⑥ 速率限制。
Q17. MQTT 3.1.1 和 5.0 的主要区别? ★★★☆☆
5.0 增加:原因码(出错知原因)、会话过期时间(替代布尔 cleanSession)、消息过期、共享订阅(原生负载均衡)、用户属性(自定义 KV 头)、主题别名(省流量)、流量控制、AUTH 报文(增强认证)。 加分点:“会话过期时间是运维上的重要改进——3.1.1 里 cleanSession=false 会产生永久僵尸会话,导致 Broker 内存泄漏。”
Q18. 消息体太大怎么办? ★★★☆☆
MQTT 协议上限 256MB,但 Broker 通常限制 1MB。 ★ 正解:大文件走对象存储,MQTT 只发 URL;或分片传输带 seq 组装。 MQTT 设计用于小报文高频,超过几十 KB 就该重新设计。
5.3.3 场景题(19~23,★★★★☆ ~ ★★★★★)
Q19. 设计一个共享单车的通信方案。 ★★★★★
答(结构化):
- 协议:MQTT(长连接、省电、服务端可推送开锁指令)。
- ClientId:
bike-{车架号},保证唯一,避免互踢。- 主题:
bike/{id}/location(上报,QoS 0 或 1)、bike/{id}/status(状态,retained)、bike/{id}/cmd(开锁指令,QoS 1)、bike/{id}/cmd/ack(回报)。- 省电:KeepAlive 设 300 秒;不上报时进入休眠,定时唤醒。
- 开锁可靠性:QoS 1 +
cmdId幂等 + 设备回报执行结果 + 超时重试 + 蓝牙兜底。- 离线检测:遗嘱 + retained 状态。
- 数据存储:EMQX 规则引擎 → Kafka → Flink 实时计算 + 时序库存轨迹。
- 安全:一车一密、TLS、ACL 限制只能访问自己的主题。 加分点:能说出“开锁要应用层回报,因为 QoS 只保证消息送达不保证锁开了”。
Q20. 设备上报的数据丢了,怎么排查? ★★★★☆
① 看发布端 QoS 是不是 0(关键数据应设 1); ② 看订阅端 QoS 是不是比发布端低(取小值降级); ③ 看设备是不是
CleanSession=true(离线消息不保留); ④ 看 ClientId 是不是固定的(随机会话找不到,离线消息补发不了); ⑤ 看 Broker 离线队列是不是满了(max_queued_messages); ⑥ 看 Broker 有没有重启。 加分点:“我们会在设备侧加’发送确认 + 本地重发’——设备本地缓存未确认的消息,收到 PUBACK 才删除,这是端到端可靠性的最后一道防线。”
Q21. 怎么实现“远程升级 10 万台设备固件”? ★★★★★
- 灰度:先 1% → 10% → 50% → 100%,每批观察成功率。
- 指令设计:
device/{id}/ota(QoS 2)带url、version、md5;设备回ota/progress报进度。- 幂等:带
otaId,设备已升级过该版本就跳过。- 校验:下载后验 MD5/SHA256,失败回滚。
- 断点续传:固件走 HTTP Range 下载,不走 MQTT。
- 进度监控:订阅
device/+/ota/progress,大屏看升级进度。- 回滚机制:失败率超阈值自动暂停,支持远程回滚到旧版本。 加分点:“★ 大文件绝不走 MQTT,MQTT 只发’去这个 URL 下载’的指令。”
Q22. 为什么不用 HTTP 轮询而用 MQTT? ★★★☆☆
① HTTP 每次带几百字节头部,MQTT 最小 2 字节,流量省几十倍; ② HTTP 每次要 TCP+TLS 握手,功耗是 MQTT 的 10 倍(电池设备致命); ③ HTTP 服务端不能主动推,做不到远程开锁; ④ HTTP 轮询有延迟,MQTT 实时; ⑤ MQTT 支持离线消息和遗嘱。 加分点:“反过来说,如果设备一天只上报一次、且不需要云端主动控制,HTTP 反而更简单——技术选型要看场景。”
Q23. MQTT 和 RabbitMQ 能互相替代吗? ★★★★☆
★ 不能,它们是不同层面的东西。
- MQTT 是设备接入协议:海量长连接、广播、不擅长堆积。
- RabbitMQ 是服务间消息中间件:队列竞争消费、事务消息、能堆积。 典型配合:设备 → MQTT → EMQX → 规则引擎 → RabbitMQ/Kafka → 业务服务。 加分点:“MQTT 是广播,所有订阅者都收到一份;RabbitMQ 队列是竞争消费,一条只被一个消费者处理——用 MQTT 做任务分发会导致所有消费者重复干活。”
5.3.4 追问链(24~25)
Q24. 追问链:可靠性
“QoS 1 会丢消息吗?” → 不会丢,但会重复 → “重复怎么办?” → 消费端幂等(msgId + SETNX + 唯一索引) → “用 QoS 2 不就不用幂等了吗?” → ★ 不能省。QoS 2 只保证消息层,业务层仍可能重复(如消费成功但写库失败后重投),幂等是消费端的铁律 → “设备执行了指令但回报丢了怎么办?” → 云端超时后查询设备当前状态(状态拉取兜底),而不是盲目重试
Q25. 追问链:运维
“Broker 内存持续上涨怎么办?” → 查僵尸会话和离线队列 → “僵尸会话哪来的?” → cleanSession=false 的设备报废后不再上线 → “怎么解决?” → MQTT 5 设 session-expiry;3.1.1 用 API 定期清理 → “设备量再涨十倍呢?” → EMQX 集群 + LB + 按地域分片 + 边缘 Broker 桥接 → “Broker 全挂了怎么办?” → ★ 设备侧要有本地缓存和重试;Broker 要集群 + 跨可用区部署;关键指令要有应用层确认
5.4 第五章小结:MQTT 核心认知
五条必须记住的认知
- ★ MQTT 是“设备接入协议”,不是“消息中间件”。它解决“百万设备连得上”,不解决“海量消息存得住”——后者交给 Kafka。
- ★ QoS 是“发布 QoS”和“订阅 QoS”取小值。想端到端可靠,两端都要设。
- ★ 消费端幂等是铁律,因为 QoS 1 会重复,QoS 2 也不能保证业务层不重复。
- ★ 遗嘱只在异常断开时触发,正常
disconnect()不触发——所以正常下线要自己发 offline。 - ★ MQTT 回调线程和心跳线程是同一个,回调里做耗时操作会导致连接被踢——必须异步化。
一页纸速查卡
【QoS】
0 最多一次(可能丢) → 高频传感器
1 至少一次(可能重复) → 状态、告警(★ 消费端幂等)
2 恰好一次(四次握手) → 计费、支付指令
⚠️ 实际投递 = min(发布QoS, 订阅QoS)
【关键机制】
保留消息 Retained → 新订阅者立刻看到最新状态;清空=发空载荷+retain
遗嘱 LWT → 异常掉线自动发;正常 disconnect 不触发
会话 CleanSession → false 才能收离线消息;ClientId 必须固定
心跳 KeepAlive → 电池设备设 300~1800 秒省电
通配符 → + 单层 / # 多层;仅订阅可用;# 不匹配 $SYS
【端口】1883 TCP | 8883 TLS | 8083 WS | 8084 WSS | 18083 EMQX控制台
【幂等三件套】msgId + Redis SETNX + 数据库唯一索引
【可靠下发】QoS1 + cmdId + 执行回报 + 超时重试 + 状态查询兜底
【大文件】走对象存储,MQTT 只发 URL
【架构】设备─MQTT─►EMQX─规则引擎─►Kafka─►业务系统
上线前检查清单
- 关闭匿名访问
allow_anonymous = false - 每个设备独立账号(一机一密),密码 bcrypt 哈希
- 配置了 ACL,设备只能访问自己的主题
- 生产启用 TLS(8883)
- ClientId 唯一且固定(需要离线消息时)
- 关键数据 QoS ≥ 1,发布端和订阅端都设
- 消费端做了幂等(msgId 去重)
- MQTT 回调里没有耗时操作(已异步化)+ try-catch
- 设置了遗嘱 + 正常下线主动发 offline
- 重连有指数退避 + 随机抖动
- 心跳间隔按设备类型设置(电池设备要长)
- 会话有过期策略(防僵尸会话)
- 限制了离线队列长度和连接速率
- 调大了 Linux
ulimit -n - 大文件走对象存储,不走 MQTT
五句话记忆法
- MQTT = 物联网的微信:Broker 是服务器,发布者是作者,订阅者是粉丝,Topic 是公众号。
- QoS 0 是平信,1 是挂号信(要签收,可能寄重复),2 是回执挂号信(来回确认四次)。
- 保留消息 = 报亭的最新样刊,遗嘱 = 登山前的遗书。
- MQTT 负责连得上,Kafka 负责存得住——上下游配合,不是替代。
- 消费端幂等是铁律,因为“至少一次”意味着“可能多次”。
配套文档:
- 名词解释 →
00-名词速查手册.md - 消息队列(RocketMQ/Kafka/RabbitMQ) →
04-中间件-讲解与面试题.md - IoT 协议安全专题 →
13-新兴攻击面与专项安全/ - 事务与一致性 →
10-数据一致性与缓存同步-实战专题.md
最后更新:2026-09-20