目录
一、先区分两个顺序
二、完整落地方案(生产级智慧充电实践)
1. 设备端:源头约束,是顺序的根基
2. 传输层选型:MQTT vs WebSocket
3. 消息中间件解耦(设备长连接网关 ↔ 订单服务)
4. 订单服务消费端:业务层保证顺序(最核心)
方案 A:单线程消费每个会话(简单,中小规模)
方案 B:单 partition,消费线程设置为 1(简单粗暴)
幂等必须配套
5. 乱序兜底:延迟消息 + 设备补报机制
6. 极端 case 处理
三、架构简图
四、避坑点(项目高频踩坑)
五、补充:如果不用 MQ,网关直接长连接推送订单服务
六、核心总结一句话
业务背景:充电桩设备通过长连接(MQTT/WebSocket)上报实时充电报文:开始充电、电压电流、功率、电量、结束充电等;订单服务消费报文生成充电流水、更新订单状态。 问题风险:网络抖动、重连、报文重传、服务异步处理,会出现时序错乱:比如先收到结束充电,再收到充电过程采样数据,导致订单状态异常、电量统计错误。
核心矛盾:网络层不保证有序 + 长连接断线重连会乱序 / 重复 + 服务端多线程消费乱序,要从「设备端、传输协议、消息标识、服务端消费、兜底补偿」五层做顺序保障。
一、先区分两个顺序
- 设备产生事件的物理时序:设备本地真实发生顺序(采样 1→采样 2→结束),这是源头,不能被破坏。
- 服务端处理顺序:服务端必须按照设备发生顺序处理,不能颠倒。
注意:长连接 MQTT/WebSocket,单 TCP 连接内网络报文本身是 TCP 有序的,但一旦发生:断线重连、设备多线程发送、QoS 重传、服务端多线程并发消费,就会出现乱序。 TCP 只保证同一个 socket 链路的字节流有序;一旦断连新建 socket,新连接消息和旧连接消息之间没有网络层顺序保证。
二、完整落地方案(生产级智慧充电实践)
1. 设备端:源头约束,是顺序的根基
- 单设备单上报队列,串行发送,禁止多线程并发上报充电报文设备内部维护充电事件 FIFO 队列,充电状态报文严格入队,一个发送完成再发下一条;不能多个线程同时往长连接写数据。
反面:设备多线程同时上报采样数据,即使同一个 TCP 连接,应用层输出顺序乱掉。
- 每条充电报文携带 3 个核心时序字段
{ "deviceId":"桩编号", "orderNo":"充电订单号", "seq":12345, // 本订单内单调递增序列号,每个充电会话从1开始,每上报一条+1 "eventTimestamp":17xxxxxx, // 设备本地事件发生时间(真实采样时间,不是发送时间) "sessionId":"桩+充电会话唯一ID", // 一次充电会话全局唯一,断线重连不变,新充电会话换新id "dataType":"sampling/start/stop" }seq:同一个充电会话内严格单调自增,不回退、不重复。设备本地内存维护,一次充电从 1 开始;断电丢失可以由设备从充电会话本地存储恢复 seq。sessionId:区分不同充电会话;旧会话消息,不允许干扰新订单。设备重启、拔枪,sessionId 变更。eventTimestamp:设备实际发生时间,用来兜底校验,不能用服务端接收时间。
- 断线重连策略 重连成功后,先补发断线期间未确认的报文,再发送新报文;不允许重连后直接发送最新数据,把历史数据丢了。
MQTT QoS1/QoS2 可以实现报文重传,但 QoS 只会保证至少一次,不保证业务顺序,业务层必须自己带 seq。
重点:TCP 有序只针对当前存活的连接。断连后旧连接滞留报文 + 新连接报文,网络到达服务端完全可能乱序,TCP 无能为力。
2. 传输层选型:MQTT vs WebSocket
智慧充电大多用 MQTT:
- MQTT 同一个 clientId,broker 内部,同一个 topic,QoS 下同一个会话内,broker 投递是有序;但是!设备断开,会话过期之后,再次重连,缓存消息和新消息,消费者多线程消费依然会乱序。
- ❗MQTT broker 只能保证 broker 内部发送顺序,不能保证消费端多线程消费顺序。很多人踩坑:以为用 MQTT 就天然有序,消费线程池并发消费直接乱序。
关键点:顺序性不能交给中间件,中间件只做投递,业务层必须做 seq 校验。
3. 消息中间件解耦(设备长连接网关 ↔ 订单服务)
架构分层:
充电桩设备 → MQTT Broker / WebSocket 网关 →消息网关服务→ RocketMQ/Kafka → 订单服务
- 网关收到设备长连接报文,不直接调用订单服务 RPC。 如果网关直接同步调用订单服务,订单服务卡顿会阻塞设备上报;同时多线程 RPC 返回无法保证顺序。
- 网关按
deviceId + sessionId作为分区 key 投递到 MQ:- Kafka:key=
deviceId_sessionId,保证同一个充电会话所有消息进入同一个 partition。同一个 partition 消息在 broker 层面是有序的。 - RocketMQ:相同 key 进入同一个队列。
- Kafka:key=
这一步非常关键:同一个充电会话全部消息落在同一个队列,避免跨队列乱序。 ⚠️ 但是:即使单 partition,如果消费者开启多线程并发消费该 partition,依然乱序!
4. 订单服务消费端:业务层保证顺序(最核心)
方案 A:单线程消费每个会话(简单,中小规模)
同一个deviceId+sessionId的消息,串行处理,内存维护当前期望序列号expectSeq。
- 收到消息:
- 如果 sessionId 已经是已结束会话(订单已完结):直接丢弃这条过期消息(历史延迟到达的采样报文)。
- 判断报文 seq:
seq == expectSeq:正常业务处理,处理完成 expectSeq +=1;seq < expectSeq:重复 / 延迟旧报文,直接幂等丢弃;seq > expectSeq:中间报文缺失,暂停消费,放入本地缓冲队列,等待缺失 seq 到达;等待超时则触发补报告警。
问题:内存缓冲,如果服务重启,内存缓冲丢失,需要持久化存储每个会话的expectSeq,存入 Redis。 Redis 存储结构:key:charging:seq:{deviceId}:{sessionId}value = 当前期望序列号,同时设置会话过期时间。
方案 B:单 partition,消费线程设置为 1(简单粗暴)
同一个 topic 下,每个设备会话落到一个 partition,消费组消费该 partition 只用1 个消费线程。
优点:代码简单,天然 broker 顺序;缺点:并发能力受 partition 数量限制,充电桩数量巨大场景不适合。
幂等必须配套
因为长连接重传,会有重复报文,seq 同时做幂等 key,避免重复扣电量、重复生成流水。
5. 乱序兜底:延迟消息 + 设备补报机制
现实网络一定会出现丢包、报文延迟(结束报文先到,采样后到):
- 当订单收到 stop 结束报文,标记订单为「待结束」,不立刻完结;开启延迟窗口(例如 30~60s)。 延迟窗口内继续接收该 sessionId 的采样报文,更新订单电量;窗口时间到,再正式闭合订单。
- 服务端检测 seq 缺口(比如收到 seq=10,但 expectSeq=7,缺 8、9),下发 MQTT 指令通知设备,重传该 session 缺失 seq 区间的历史充电数据。
- 会话超时:超过最大充电时长,自动关闭会话,清理 Redis 中 seq 状态。
6. 极端 case 处理
- 结束报文先到达,采样报文后到达session 收到 stop 报文,进入延迟窗口期;后续迟到的采样 seq 只要合法,依旧更新流水;窗口期结束,直接忽略迟到报文。
- 设备重启,sessionId 更新旧 sessionId 消息全部做过期丢弃;新 session 从 seq=1 重新开始。新旧会话完全隔离,互不干扰。
- 服务实例重启所有会话的
expectSeq持久化在 Redis,重启后读取 Redis 继续校验 seq,不会丢失顺序状态。 - 重复重传报文(MQTT QoS1 至少一次)
seq < expectSeq直接丢弃,天然幂等。
三、架构简图
充电桩设备(本地FIFO队列,携带seq+sessionId) ↓长连接MQTT MQTT Broker ↓网关转发,key=device+sessionId RocketMQ/Kafka(同会话进入同一个队列partition) ↓ 订单服务消费 ├─ Redis维护每个会话expectSeq ├─ seq校验、缺口缓冲 ├─ stop报文开启延迟窗口 └─ 缺口下发补报指令给设备四、避坑点(项目高频踩坑)
- ❌ 只依赖 TCP/MQTT 保证顺序,业务不加 seq。断连重连后乱序,订单电量错乱。
- ❌ Kafka 单 partition,但是消费者配置多线程消费同一个 partition。partition 有序,多线程消费之后处理完全乱序。
- ❌ 使用服务端接收时间做排序。网络延迟,接收时间不能代表真实充电发生顺序。必须以设备上报 eventTimestamp+seq 为主。
- ❌ 内存保存 expectSeq,服务重启状态丢失,顺序校验失效,必须 Redis 持久化。
- ❌ 收到 stop 报文直接关闭订单,后面迟到采样报文无法更新电量,统计少算 / 多算电量。
五、补充:如果不用 MQ,网关直接长连接推送订单服务
如果架构是 WebSocket/MQTT 网关直接调用订单服务(不走消息队列): 网关层需要做会话级串行:同一个deviceId+sessionId,网关内部维护内存队列,串行推送给订单服务,上一个业务 ACK 返回,再推下一条。 一旦网关集群部署,设备会漂移到不同网关实例,内存队列失效,所以这种架构不适合集群生产环境,强烈建议引入消息中间件做会话分区。
六、核心总结一句话
TCP/MQTT 只能保证单连接内网络投递有序;断线重连、集群、多线程消费都会打破顺序。真正保证充电业务顺序:设备端单调 seq+sessionId,消息按会话分区投递,消费端持久化维护期望序列号,stop 结束做延迟窗口兜底,缺口触发设备补报。