QoS2 消息在 10 万 TPS 下出现重复,如何验证 Exactly-Once 语义
解读
- 场景定位:MQTT QoS2 在国内 IoV、移动支付、IM 等高并发链路中普遍使用,10 万 TPS 已接近头部互联网公司的峰值,重复消费会直接造成订单双花、计费和库存差错,业务零容忍。
- 问题本质:QoS2 协议理论上提供“Exactly-Once”,但在 Broker 实现、客户端重传、网络闪断、GC 停顿、负载均衡漂移等情况下,仍可能出现重复投递。性能测试工程师的任务不是“跑高并发”,而是“在高并发下用可量化的方法证明没有重复”。
- 面试考点:
- 是否理解 MQTT 5.0 QoS2 四步握手、PacketId 生命周期、Broker 去重表机制;
- 能否把“去重”抽象成可观测、可度量、可灰度的验证模型;
- 是否具备在 10 万 TPS 规模下设计低成本、低干扰、高可信的测试方案的能力;
- 是否熟悉国内常用中间件(EMQX、HiveMQ、阿里云 IoT Hub、腾讯云 IoT)的重复触发场景。
知识点
- MQTT QoS2 四步握手:PUBLISH → PUBREC → PUBREL → PUBCOMP,PacketId 在完成 PUBCOMP 前必须被 Broker 和 Client 视为“in-flight”。
- 重复根因:
- Broker 侧:PUBREC 丢失后 Client 重传 PUBLISH,Broker 未幂等;
- Client 侧:PUBCOMP 丢失后 Broker 重传 PUBREL,Client 重复回调;
- 水平扩容:Sticky 失效导致不同节点各自维护 PacketId 空间;
- 持久化:去重表刷盘延迟,Broker 重启恢复时丢失状态。
- 验证模型:
- 端到端唯一标识:业务层 Snowflake/UUID + MQTT 层 PacketId 双因子;
- 原子计数:每个消息在 Broker In-flight 表、Consumer 回调、业务 DB 三个点各打一次“指纹”,最终三者总数差为 0;
- 时间窗口去重:Client 侧维护 5 秒滑动窗口,窗口内 PacketId+Topic 重复即丢弃;
- 对账:测试结束后用 Flink Batch 对“生产 UUID 集合”与“消费 UUID 集合”做差集,差集必须为空。
- 性能测试工具:
- JMeter 5.5 + MQTT plugin 仅支持 1 万 TPS 级,需二次开发 Netty 连接池;
- 国内主流自研压测平台:阿里 PTS、腾讯 WeTest、字节 BitSail,支持 100 万并发长连接;
- 观测:EMQX 内置 Prometheus exporter,重点看
emqx_messages_qos2_received、emqx_messages_qos2_dropped、emqx_messages_qos2_duplicate。
- 资源预算:10 万 TPS QoS2 平均 256 Byte 报文,网络吞吐约 25 MB/s,Broker 16 核 32 GB 三节点即可,但需额外 200 GB NVMe 存放去重表,避免磁盘 IO 成为瓶颈。
答案
步骤 1:建立“零重复”基线
- 在 1 千 TPS 低压下运行 30 min,确认 Flink 差集为空,作为“无重复”基线。
步骤 2:构造 10 万 TPS 稳态
- 使用自研 Netty 客户端,单机 5 万连接、每连接 2 msg/s,共 4 台压测机;
- 消息体带 64 位 Snowflake ID、时间戳、压测机 IP,方便对账;
- 开启 TCP_NODELAY、SO_KEEPALIVE=30 s,避免因小包粘包导致 PacketId 错乱。
步骤 3:注入真实故障
- 网络:tc 命令 1% 丢包、200 ms 延迟 10 s 周期,模拟移动网络;
- Broker:kill -9 随机节点,验证集群重平衡后去重表是否漂移;
- Client:Full GC 暂停 3 s,触发重传。
步骤 4:实时断言
- Grafana 看板:
- 曲线 A:
emqx_messages_qos2_received累计值; - 曲线 B:Flink 消费算子输出的“首次出现 UUID” 累计值;
- 曲线 C:业务 DB 落库“INSERT COUNT” 累计值;
三条曲线必须完全重合,差异 >0 即告警停压。
- 曲线 A:
步骤 5:离线对账
- 压测结束后,将生产 UUID 文件、消费 UUID 文件、DB 主键分别排序后做
diff -u; - 若差集非空,用 PacketId+Topic+时间戳 回溯日志,定位是 Broker 还是 Client 的重复;
- 重复率 = 重复条数 / 总条数,目标 <1 ppm(百万分之一)。
步骤 6:出具报告
- 给出“在 10 万 TPS、30 min、3 轮故障注入下,重复率为 0 ppm,Exactly-Once 语义成立”结论;
- 附关键配置:Broker
max_inflight=200、session.replayq=true、Clientclean_start=false、session_expiry_interval=300 s; - 风险提示:当 TPS >15 万或消息 >1 KB 时,磁盘去重表 RTT 升高,需升级为内存+Redis 共享方案。
拓展思考
- 成本优化:若业务可接受“至少一次”+幂等消费,是否值得为 QoS2 付出 30% 的 RT 损耗?可用 A/B 实验对比 QoS1+业务幂等 与 QoS2 的资源成本。
- Serverless 场景:函数计算冷启动 1 s,会导致 PUBREC 超时重传,Broker 去重表需延长到 10 s,可能撑爆内存,如何动态伸缩?
- 多协议混用:同一主题下既有 MQTT QoS2 也有 HTTP 回调,如何统一“Exactly-Once”语义?可引入 Kafka 事务消息作为唯一日志,所有协议都向该日志幂等写。
- 合规审计:国内金融客户要求“每笔消息可追踪”,需把 PacketId、UUID、TraceId 落盘 5 年,性能测试方案要同步验证 50 TB 级冷数据查询在 30 s 内返回,避免过了性能关却过不了审计关。