隧道里断网不到一分钟,工单系统却在服务端多出三条“已到场”。设备端日志看起来很正常:QoS 1 发布成功、重连完成、PUBACK 也回来了。真正的问题不在有没有确认,而在断线前后同一个业务动作被重新装进不同的 MQTT 报文,packetId 变了,服务端就把它们当成两次操作。

这次我把 Demo 做成 FieldOutboxLab,页面是 SyncLedgerPage,任务号 OUTBOX-1408。稳定状态固定为:客户端 field-node-07 已连接,sessionPresent=true,QoS 1,最近 packetId 184,出站箱待发 2 条、飞行中 1 条、已确认 37 条,服务端丢弃重复 3 条,重连 2 次,最近确认耗时 86ms,最终状态 SYNC_STABLE。

一、QoS 1 解决送达,不替业务判断“是不是同一单”

最早的实现把 packetId 当成去重键。它在单次连接里看起来可用:消息发布后收到 PUBACK,就把对应记录改成完成。问题是 packetId 属于 MQTT 协议会话,有自己的分配与复用规则;应用重启、会话重建或发送窗口滚动后,同一个数字可能再次出现。反过来,同一个业务命令在重新排队时也可能拿到另一个 packetId。

因此我把两类标识拆开。commandId 在用户点击时生成,代表“工单 WO-20261003-087 的到场动作”,进入 RDB 后不再改变;packetId 只作为协议诊断字段,记录最后一次发送尝试。服务端按 commandId 做幂等,客户端按 commandId 更新出站箱。即使 packetId 从 176 变成 184,业务仍知道这是同一条命令。

这与上一轮的 WebSocket 序列补洞也不是同一个问题。行情流关心连续序列与快照水位,允许丢弃过时更新;工单命令是离散副作用,不能靠最新快照覆盖。一个“确认到场”重复执行,可能触发计费、通知和审计,所以必须保留稳定业务身份。

二、出站箱先落库,再允许网络层看见命令

第一段代码解决进程在点击后、发送前退出的问题。enqueueCommand() 在 RDB 事务中写入 commandId、topic、payload、状态和尝试次数;提交成功后才通知发送器。页面上的“待发 2”来自数据库查询,不是内存数组长度。

type OutboxState = 'PENDING' | 'INFLIGHT' | 'ACKED'

interface OutboxCommand {
  commandId: string
  topic: string
  payload: string
  state: OutboxState
  attempt: number
  lastPacketId: number
  createdAt: number
}

async function enqueueCommand(store: relationalStore.RdbStore,
  workOrderId: string): Promise<string> {
  const commandId = `ARRIVE-${workOrderId}-${Date.now()}`
  await store.insert('outbox_command', {
    command_id: commandId,
    topic: 'field/workorder/arrive',
    payload: JSON.stringify({ commandId, workOrderId, action: 'ARRIVE' }),
    state: 'PENDING', attempt: 0, last_packet_id: 0, created_at: Date.now()
  })
  return commandId
}

示例省略了表结构和事务包装,实际项目会把业务状态变更与 outbox 写入放进同一事务,避免页面显示“已到场”而命令没有入箱。commandId 不能只用时间戳,产品代码还组合设备标识和随机分量;本文使用可读格式是为了让日志与图片更容易核对。

重复点击的风险也在入库前处理。同一工单、同一动作存在未完成记录时,按钮不会新建 commandId,而是回显当前任务。页面离场不删除 PENDING 项,因为它属于业务队列,不属于页面生命周期;监听器可以释放,队列数据必须留下来供下次启动继续发送。

三、重连时先判断会话,再决定如何重放

第二段代码把 mqtt.js 连接事件收敛成明确状态。客户端使用稳定 clientId field-node-07,关闭 clean session;连接回调记录 sessionPresent。无论 broker 是否保留会话,业务层都从 RDB 恢复 PENDING 与超时 INFLIGHT,但不会凭 packetId 直接认定已确认。

function connectBroker(): MqttClient {
  const client = mqtt.connect(BROKER_URL, {
    clientId: 'field-node-07',
    clean: false,
    reconnectPeriod: 2000,
    connectTimeout: 8000
  })

  client.on('connect', async (packet) => {
    ledger.sessionPresent = packet.sessionPresent === true
    ledger.reconnects++
    await recoverTimedOutInflight()
    await flushOutbox(client)
  })

  client.on('packetreceive', (packet) => {
    if (packet.cmd === 'puback') ledger.lastPacketId = packet.messageId ?? 0
  })
  client.on('offline', () => ledger.setState('WAIT_NETWORK'))
  return client
}

sessionPresent=true 说明 broker 识别了既有会话,不代表业务命令已经落到服务端数据库。PUBACK 也只证明协议侧确认,服务端仍可能在业务提交前失败。因此本项目的响应主题会带回 commandId 和业务结果,客户端收到业务 ACK 后才把记录改成 ACKED;PUBACK 仅更新诊断字段和超时计时。

重复注册事件是另一个隐患。页面每次出现都创建 client,会让同一条 ACK 被多个监听器处理。现在连接由应用级 BrokerSession 持有,页面只订阅状态快照;应用退后台时不销毁出站箱,是否保持网络连接则由产品策略决定。真正关闭 session 时必须成对移除监听并调用 end()。

四、重放要有抢占状态,不能边查边发

第三段代码处理并发重放。flushOutbox() 先用条件更新把一条 PENDING 抢成 INFLIGHT,只有更新行数为 1 的执行者可以发布。publish payload 始终包含 commandId,发送失败则回滚为 PENDING,并增加 attempt;超过上限进入人工可见的暂停状态,而不是无限重试。

async function flushOutbox(client: MqttClient): Promise<void> {
  if (!client.connected || ledger.flushing) return
  ledger.flushing = true
  try {
    for (const cmd of await outboxRepo.listReplayable(20)) {
      const claimed = await outboxRepo.claim(cmd.commandId, 'PENDING', 'INFLIGHT')
      if (!claimed) continue
      await new Promise<void>((resolve, reject) => {
        client.publish(cmd.topic, cmd.payload, { qos: 1, retain: false },
          (error?: Error) => error ? reject(error) : resolve())
      }).catch(async () => {
        await outboxRepo.returnPending(cmd.commandId, cmd.attempt + 1)
      })
    }
  } finally {
    ledger.flushing = false
  }
}

这段 callback 表示 mqtt.js 完成相应发送流程,不替代业务 ACK。客户端收到服务端响应后,以 commandId 执行 markAcked();如果同一响应到达两次,第二次更新行数为 0,只增加 duplicateDropped 诊断,不重复推动页面。RDB ResultSet 使用完立即 close,发送器停止时取消定时器,避免后台残留轮询。

抢占状态还有崩溃边界。进程可能在改成 INFLIGHT 后立即退出,所以记录保存 claimedAt。下次启动只回收超过 15 秒的 INFLIGHT,不碰仍可能收到 ACK 的新记录。这个超时不是网络协议常量,而是 Demo 的产品参数,弱网环境需要结合实际链路调整。

五、日志要同时呈现协议层和业务层

工程目录分为 OutboxRepository.ets、BrokerSession.ets、ReplayCoordinator.ets、SyncLedger.ets 和 SyncLedgerPage.ets。HiLog 固定输出:task=OUTBOX-1408 clientId=field-node-07 connected=true sessionPresent=true、qos=1 lastPacketId=184 reconnects=2 latency=86ms、pending=2 inflight=1 acked=37 duplicateDropped=3、state=SYNC_STABLE。

IDE 右侧模拟器展示与日志相同的账本。项目树、代码和手机状态并不是三套各自编写的文案:SyncLedger 是单一数据源,页面只渲染快照。红色标注指向 commandId 去重,不去圈 packetId,因为后者恰恰不能承担业务身份。

六、这次恢复测试从“断在哪”开始

我把网络中断放在三个位置:命令入库前、PUBLISH 发出后但 PUBACK 前、PUBACK 后但业务 ACK 前。第一种不会产生队列记录;第二种重连后允许协议重发,服务端用 commandId 去重;第三种客户端会再次查询或重放,服务端仍返回同一业务结果。三条路径最终都只形成一条到场记录。

最终页时间 14:08、电量 82%,任务 OUTBOX-1408。当前仍有待发 2、飞行中 1,不代表失败,而是模拟弱网队列的运行中快照;已确认 37、重复丢弃 3、重连 2、最近耗时 86ms,状态 SYNC_STABLE。按钮“模拟断网”制造网络切换,“重放出站箱”触发一次受控 flush。

七、协议可靠不等于产品可靠

QoS 1 很有价值,它避免应用自己重造全部传输确认;可一旦命令具有业务副作用,就不能把协议确认当成最终事实。packetId、commandId、PUBACK 和业务 ACK 各自回答不同问题:报文是哪一份、动作是哪一次、协议是否收到、业务是否提交。把它们压成一个字段,弱网下迟早会出现解释不清的重复。

出站箱也不是永远增长的日志。ACKED 记录按审计周期归档,失败超过次数的记录进入可见队列,用户注销时按业务规则清理。清理任务与发送任务使用同一把仓库锁,防止刚确认的记录被错误重放。应用卸载后的服务器幂等窗口也要独立设计,不能依赖设备本地仍保存 commandId。

这次改造最后没有增加更多“自动重试”按钮,而是把每条命令的状态、尝试和身份都做成可观察账本。实际用下来,弱网同步最难的从来不是多发一次请求,而是系统能不能证明多发的那一次仍然是同一件事。

Logo

讨论HarmonyOS开发技术,专注于API与组件、DevEco Studio、测试、元服务和应用上架分发等。

更多推荐