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 是为谁设计的、和 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 的存在

关键点:

  1. 发布者只管往主题发,不关心谁订阅。
  2. 订阅者只管订阅主题,不关心谁发的。
  3. Broker 负责匹配和转发。
  4. 一个客户端既可以是发布者也可以是订阅者(比如一个智能灯:订阅“开灯指令”,发布“当前状态”)。

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) ──────────────│

★ 关键理解:“至少一次”意味着消息一定到达,但可能到达多次。

为什么可能重复:

  1. 发送方发出 PUBLISH 后,在收到 PUBACK 之前连接断了。
  2. 发送方重连后,不确定 Broker 到底收到没有,只能重发(DUP=1)。
  3. 如果 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 为客户端保留的“档案”,包含:

  1. 客户端的订阅关系
  2. 已发送但未确认的消息(QoS 1/2)
  3. 已收到但未完成的消息(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

机制:

  1. 客户端 CONNECT 时声明 KeepAlive 秒数(如 60)。
  2. 客户端必须在 1.5 × KeepAlive 时间内至少发一个报文(可以是 PINGREQ,也可以是任意 PUBLISH)。
  3. 如果客户端没数据要发,就发 PINGREQ,Broker 回 PINGRESP。
  4. 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);
        }
    }
}

★ 这个案例体现的三个关键设计:

  1. 遗嘱 + 主动上报组合,准确感知设备上下线。
  2. 指令 ID(cmdId) 做幂等,杜绝重复执行。
  3. 应用层回报机制弥补 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. 设计一个共享单车的通信方案。 ★★★★★

答(结构化):

  1. 协议:MQTT(长连接、省电、服务端可推送开锁指令)。
  2. ClientId:bike-{车架号},保证唯一,避免互踢。
  3. 主题:bike/{id}/location(上报,QoS 0 或 1)、bike/{id}/status(状态,retained)、bike/{id}/cmd(开锁指令,QoS 1)、bike/{id}/cmd/ack(回报)。
  4. 省电:KeepAlive 设 300 秒;不上报时进入休眠,定时唤醒。
  5. 开锁可靠性:QoS 1 + cmdId 幂等 + 设备回报执行结果 + 超时重试 + 蓝牙兜底。
  6. 离线检测:遗嘱 + retained 状态。
  7. 数据存储:EMQX 规则引擎 → Kafka → Flink 实时计算 + 时序库存轨迹。
  8. 安全:一车一密、TLS、ACL 限制只能访问自己的主题。 加分点:能说出“开锁要应用层回报,因为 QoS 只保证消息送达不保证锁开了”。

Q20. 设备上报的数据丢了,怎么排查? ★★★★☆

① 看发布端 QoS 是不是 0(关键数据应设 1); ② 看订阅端 QoS 是不是比发布端低(取小值降级); ③ 看设备是不是 CleanSession=true(离线消息不保留); ④ 看 ClientId 是不是固定的(随机会话找不到,离线消息补发不了); ⑤ 看 Broker 离线队列是不是满了(max_queued_messages); ⑥ 看 Broker 有没有重启。 加分点:“我们会在设备侧加’发送确认 + 本地重发’——设备本地缓存未确认的消息,收到 PUBACK 才删除,这是端到端可靠性的最后一道防线。”

Q21. 怎么实现“远程升级 10 万台设备固件”? ★★★★★

  1. 灰度:先 1% → 10% → 50% → 100%,每批观察成功率。
  2. 指令设计:device/{id}/ota(QoS 2)带 url、version、md5;设备回 ota/progress 报进度。
  3. 幂等:带 otaId,设备已升级过该版本就跳过。
  4. 校验:下载后验 MD5/SHA256,失败回滚。
  5. 断点续传:固件走 HTTP Range 下载,不走 MQTT。
  6. 进度监控:订阅 device/+/ota/progress,大屏看升级进度。
  7. 回滚机制:失败率超阈值自动暂停,支持远程回滚到旧版本。 加分点:“★ 大文件绝不走 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 核心认知

五条必须记住的认知

  1. ★ MQTT 是“设备接入协议”,不是“消息中间件”。它解决“百万设备连得上”,不解决“海量消息存得住”——后者交给 Kafka。
  2. ★ QoS 是“发布 QoS”和“订阅 QoS”取小值。想端到端可靠,两端都要设。
  3. ★ 消费端幂等是铁律,因为 QoS 1 会重复,QoS 2 也不能保证业务层不重复。
  4. ★ 遗嘱只在异常断开时触发,正常 disconnect() 不触发——所以正常下线要自己发 offline。
  5. ★ 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

五句话记忆法

  1. MQTT = 物联网的微信:Broker 是服务器,发布者是作者,订阅者是粉丝,Topic 是公众号。
  2. QoS 0 是平信,1 是挂号信(要签收,可能寄重复),2 是回执挂号信(来回确认四次)。
  3. 保留消息 = 报亭的最新样刊,遗嘱 = 登山前的遗书。
  4. MQTT 负责连得上,Kafka 负责存得住——上下游配合,不是替代。
  5. 消费端幂等是铁律,因为“至少一次”意味着“可能多次”。

配套文档:

最后更新:2026-09-20