MQTT 消息可靠性实战:丢失、重复、乱序、积压的系统性治理
做物联网后端这几年,我处理过的绝大多数「线上疑难杂症」最终都能归因到四类消息问题上:丢了、重了、乱了、堆了。这篇文章把每一类的成因和治理方案系统性地整理出来——都是真实场景里踩过、验证过的。
一、先搞清楚 QoS 的真实语义
MQTT 的三个 QoS 等级,语义必须掰开揉碎理解:
| QoS | 语义 | 交互 | 真实代价 |
|---|---|---|---|
| 0 | 至多一次(At most once) | 单向发 | 丢了就丢了,TCP 层保证不了应用层不丢 |
| 1 | 至少一次(At least once) | PUBACK 二次 | 必然可能重复,幂等是消费端义务 |
| 2 | 恰好一次(Exactly once) | 四次握手 | 开销大、吞吐低,海量设备场景基本不划算 |
实战选型的基本盘:
- 控制指令(下发开灯/重启):QoS1 + 幂等设计。丢了不行,重复无害(幂等兜住)
- 遥测数据(温度/电流上报):QoS1 允许偶尔重复;高频低价值数据可以直接 QoS0
- 计费/计量类数据:唯一值得考虑 QoS2 的场景,但更常见的做法还是 QoS1 + 去重表
经验:不要迷信 QoS2。真实系统里「QoS1 + 业务幂等」的吞吐和稳定性远好于 QoS2,而且语义边界更清楚。
二、消息为什么会丢
按链路顺序排查,丢消息的场景就这几处:
设备 ──①──> Broker ──②──> 规则引擎/消费端 ──③──> 存储
弱网/TCP半开 Broker重启/积压淘汰 消费ACK前崩溃- 设备→Broker:设备断电瞬间,TCP 半开连接上 broker 还以为链路活着,QoS1 消息发出去 PUBACK 没回来,设备重连后重发——这其实是「丢」的对立面(重),但很多人在这一步误判
- Broker 层:非持久会话(cleanSession=true)断线期间的订阅消息直接丢弃;Broker 积压触发 max queued messages 淘汰策略,默认淘汰旧消息——这是最隐蔽的丢法
- 消费端:拉到消息、处理一半崩溃、ACK 没提交,重平衡后消息会重投(又是重);但如果用的是「先 ACK 后处理」的错误姿势,崩了就是真丢
对应治理:
# Broker 侧(EMQX 为例)
mqtt.max_qos = 1
mqtt.session.upgrade_qos = false # 不偷偷升级 QoS
mqtt.mqueue.max_length = 10000 # 积压上限要显式评估
mqtt.mqueue.default_priority = high # 控制类消息优先// 消费端:处理完再 ACK,是消息不丢的最后防线
consumer.subscribe(topic);
while (running) {
Message msg = consumer.poll();
try {
handle(msg); // 业务处理
consumer.commitSync(); // 处理成功才提交位点
} catch (Exception e) {
sendToRetryQueue(msg); // 失败进重试队列,不无限本地重投
}
}三、重复消息:幂等是消费端的义务
QoS1 的重复不可避免(网络抖动重发、消费端重平衡重投),所以接收侧必须幂等。三层防线:
1. 传输层去重(挡住大部分)
设备端每条消息带 client 内单调递增 seq,云端用 Redis 做滑动去重:
def is_duplicate(device_id, seq):
key = f"dedup:{device_id}:{seq}"
return not redis.set(key, 1, nx=True, ex=600) # 10 分钟窗口注意 TTL 必须大于设备断网补传的最大时间窗,否则补传的旧消息会被当成新消息。
2. 业务层幂等(最终兜底)
指令类消息按业务键加唯一约束,数据库天然挡重复:
INSERT INTO device_command (cmd_id, device_id, action, created_at)
VALUES (?, ?, ?, ?);
-- cmd_id 唯一索引,重复插入直接失败 = 幂等生效3. 状态机幂等(防错乱)
有些重复消息不该简单忽略,而是要按状态机判断是否还有效:
if (device.getState() != State.OFFLINE) {
return; // "上线消息"只有设备处于 OFFLINE 时才生效
}四、乱序:QoS 不保证跨重连的顺序
这是最容易被忽略的一类问题:同一条 TCP 连接内 QoS1 基本有序,但重连之后顺序全乱。典型事故:设备断网期间本地缓存了 3 条消息(温度 20→25→30),补传时因为重试节奏不同到达顺序变成 30→20→25,云端「最新值」显示 20。
治理方案按数据特性选:
- 最新值语义(状态上报):消息带设备端时间戳 + 单调 seq,云端 last-write-wins 按
max(seq)落库,不按到达顺序 - 事件流语义(开关记录):云端按 seq 开窗口缓存重排,窗口满了或超时了才落库
- 聚合语义(统计数据):干脆不等顺序,全部收齐后按时间桶聚合
// 状态上报的乱序防御:旧消息直接丢弃
if (incoming.seq <= device.getLastSeq()) {
metrics.staleDropped(); // 打点观测,乱序率是链路健康度指标
return;
}
device.setLastSeq(incoming.seq);五、离线重连:抖动风暴是生产事故的头号来源
几万台设备挂在同一个机房,一旦网络抖动或 Broker 重启,所有设备同时重连——连接风暴直接把 broker 打趴,然后更多设备掉线,雪崩。
设备端:指数退避 + 随机抖动
// 设备端重连策略:退避 + 抖动,错开重连时刻
delay = min(BASE * 2^attempt, MAX) // 1s, 2s, 4s, ... 上限 5min
delay = delay * (0.5 + random()) // 抖动因子,打散雪崩会话保持:重连不重建
- MQTT 3.1.1:
cleanSession=false,broker 保留会话与离线消息队列 - MQTT 5.0:
session expiry interval显式声明会话存活时长,更精细 - 设备端要有固定的
ClientID规则(如{productKey}-{deviceId}),ClientID 变了等于换了个人,会话和订阅全部作废
云端:梯度放行
Broker 侧配置连接速率限制,超出的连接快速失败(设备会退避重试,比排队挂着好)——这是「宁可让设备多等几秒,也不能让 broker 被打死」的取舍。
六、消息积压与流量削峰
设备侧整点上报、批量唤醒、固件升级回传——都会造成瞬间洪峰。三层削峰:
- 设备侧削峰(最有效):变化上报 + 周期全量。温度没变就别报,每 5 分钟报一次心跳值即可,能砍掉 90%+ 的无效上报
- 接入层缓冲:broker → Kafka 的桥接层天然削峰,Kafka 分区数按峰值吞吐 × 3 规划
- 消费层扩容:监控 consumer lag,积压超过阈值自动扩容消费者(分区数是上限)
积压告警别只看 lag 绝对值,看消费延迟(最新消息时间 - 消费位点时间)——lag 10 万条可能是 2 秒的事,也可能是 2 小时的事。
七、死信与重试:失败消息的归宿
处理失败的消息不能无限重试(会阻塞正常消息),标准路径:
正常队列 ──失败──> 重试队列(1min) ──失败──> 重试队列(5min) ──失败──> 重试队列(30min) ──失败──> 死信队列
│
告警 + 人工介入 ←──────────┘死信队列的运营纪律:
- 死信必须有告警,不能让死信队列变成「消息坟场」
- 死信要带完整的失败上下文(异常栈、重试次数、设备信息),否则人工排查无从下手
- 提供一键重放工具:修复 bug 后批量重新投递
- 每周看死信 TOP 原因分布,往往是设备固件 bug 的最早信号
八、一张排查速查表
| 症状 | 首先查 | 常见根因 |
|---|---|---|
| 状态显示开着,实际是关的 | 指令是否真正下发(broker 日志) | QoS0 丢了 / 消费先 ACK 后处理 |
| 同一指令执行了两次 | 去重表命中率 | 消费重平衡重投,幂等没做 |
| 曲线出现"回到过去"的毛刺 | 设备 seq 连续性 | 补传乱序,没按 seq 丢弃旧值 |
| 大面积设备集体掉线 | broker 连接数曲线 | 连接风暴雪崩,没做退避 |
| 页面数据延迟越来越大 | consumer lag / 消费延迟 | 消费能力不足,积压 |
| 部分设备数据永久缺失 | broker mqueue 淘汰日志 | 积压超限淘汰,默认丢旧消息 |
写在最后
物联网消息可靠性的本质不是某个组件多强,而是每一跳都明确「丢失/重复/乱序各自的兜底策略」:设备端管好缓存与补传,broker 管好会话与积压策略,消费端管好幂等与位点。把这三段的责任边界画清楚,90% 的「玄学问题」都会变成可观测、可定位的工程问题。