【WebSocket+MQTT】day61-backfill

作者:mario 发布时间: 2026-09-03 阅读量:5 评论数:0

Day 61 · 断线重连与数据补传 — 让丢掉的时间自己长回来

本周到目前为止的链路有个致命缺陷:断线 30 秒重连成功后,这 30 秒的数据永久消失了——曲线上一段刺眼的空白,值班员无从判断空白里有没有发生过超温。今天解决实时系统最难也最能体现工程功力的问题:数据补传(backfill)。涉及三个子问题:客户端怎么知道该补哪段(游标)、服务器怎么存历史(环形缓冲)、补传数据怎么和实时流无缝衔接(去重与排序)。昨天 MQTT QoS 1 的"至少一次 + 业务去重"思想,今天全部落地。


目录


一、问题定义:补传的三个子问题

时间线:
────────●─────────────✂─────────────●────────▶
      t0 断线          (30s 黑洞)      t1 重连成功

子问题 1(知道补哪段):客户端断线前收到的最后一条数据的时间戳 lastTs —— "游标"
子问题 2(历史从哪来):服务器必须留存历史数据,按游标查出缺口区间
子问题 3(如何衔接)  :补传数据 + 实时数据可能重叠/乱序 → 去重 + 排序

1.1 与 MQTT 会话保持的对照

机制 MQTT(Broker 内建) 自建 WS(今天自己写)
断线暂存 clean:false + QoS1 会话 服务器环形缓冲(第三节)
重连恢复 自动重投会话消息 客户端发 backfill 请求(第四节)
去重 业务按消息 id 去重 时间戳去重(第五节)

二、协议扩展:游标与补传消息

在 Day 57 的 WsMessage 上扩展三个消息类型:

/**
 * 协议消息全集(本周最终版)
 */
export type WsMessage =
  // 实时数据(服务器 → 客户端,每秒)
  | { type: "snapshot"; payload: PanelSnapshot; ts: number }
  // 服务器在 snapshot 里附带的最新序号(客户端据此更新游标)
  | { type: "snapshot"; seq: number; payload: PanelSnapshot }
  // 补传请求(客户端 → 服务器,重连成功后发)
  | { type: "backfill"; fromSeq: number }
  // 补传响应(服务器 → 客户端,一批历史数据 + 是否还有更多)
  | { type: "backfill-data"; fromSeq: number; snapshots: PanelSnapshot[]; hasMore: boolean }
  // 心跳
  | { type: "ping"; seq: number }
  | { type: "pong"; seq: number };

2.1 为什么用序号(seq)而不是时间戳当游标

时间戳游标的问题:
- 服务器和客户端时钟不同步 → 边界判断误差
- 同一毫秒多条消息 → 边界重复或遗漏

序号游标:
- 服务器单调整数,天然严格全序,无歧义
- 时间戳仍保留在 payload 里(显示用),游标职责分离

🎯 “用单调递增 ID 做游标,时间戳只做显示” 是分布式系统的通用律条(Kafka 的 offset、MySQL 的 binlog position 全是这个思想)——今天第一次亲手实践它。


三、服务器端:环形历史缓冲

3.1 设计

环形缓冲(定长数组,写满覆盖最旧):
容量 600 条(10 分钟 @ 1条/秒)

push 序列: n-3, n-2, n-1, n, [覆盖最旧] ...
读指针逻辑:backfill(fromSeq) → 从缓冲里找 seq >= fromSeq 的连续段

3.2 实现

/**
 * 定长环形缓冲:服务器留存最近 N 条快照
 * 写满覆盖最旧——内存恒定,历史窗口有限(10 分钟内可补)
 */
export class RingBuffer<T extends { seq: number }> {
  private items: T[] = [];
  constructor(private capacity: number) {}

  /** 追加(自动覆盖最旧,同时维护 seq 单调) */
  push(item: T): void {
    this.items.push(item);
    if (this.items.length > this.capacity) {
      this.items.shift();   // 容量小(百级)时 shift 足够;万级以上用首尾指针法
    }
  }

  /**
   * 查询 seq 大于 fromSeq 的连续段(补传核心)
   * @returns 数据 + 起始游标是否还在缓冲内(太旧的缺口无法补全)
   */
  queryAfter(fromSeq: number, limit: number): { items: T[]; hasMore: boolean; gapLost: boolean } {
    const idx = this.items.findIndex((it) => it.seq > fromSeq);
    if (idx === -1) {
      // 两种情况合并处理:请求的游标太新(没有数据要补)或太旧(缓冲头之后)
      return { items: [], hasMore: false, gapLost: fromSeq < (this.items[0]?.seq ?? 0) - 1 };
    }
    const items = this.items.slice(idx, idx + limit);
    const hasMore = idx + limit < this.items.length;
    return { items, hasMore, gapLost: false };
  }

  /** 缓冲内最旧序号(客户端可据此判断缺口是否可补) */
  get oldestSeq(): number {
    return this.items[0]?.seq ?? 0;
  }
}

3.3 服务器消息处理扩展

import { WebSocketServer, WebSocket } from "ws";
import { RingBuffer } from "./ring-buffer";

const wss = new WebSocketServer({ port: 8080, path: "/sensor-feed" });
const history = new RingBuffer<{ seq: number; payload: PanelSnapshot }>(600);
let seq = 0;

wss.on("connection", (ws: WebSocket) => {
  // 实时推送(每秒)—— 同时写入历史缓冲
  const timer = setInterval(() => {
    if (ws.readyState !== WebSocket.OPEN) return;
    const payload = generateSnapshot();          // 模拟数据生成
    const item = { seq: ++seq, payload };
    history.push(item);
    ws.send(JSON.stringify({ type: "snapshot", seq: item.seq, payload }));
  }, 1000);

  ws.on("message", (raw) => {
    const msg = JSON.parse(raw.toString());

    if (msg.type === "backfill") {
      // 补传:按游标回放历史(分批,每批最多 120 条防大帧卡连接)
      const { items, hasMore, gapLost } = history.queryAfter(msg.fromSeq, 120);
      ws.send(JSON.stringify({
        type: "backfill-data",
        fromSeq: msg.fromSeq,
        snapshots: items,
        hasMore,
        gapLost,     // true 时客户端应提示"部分历史已不可恢复"
      }));
    }
  });

  ws.on("close", () => clearInterval(timer));
});

四、客户端:会话恢复流程

4.1 状态图

        ┌──────────────────────────────────────────┐
        │ open ──正常收 snapshot──▶ 更新游标 lastSeq  │
        │   │                                       │
        │   │ 断线(心跳判死/1006)                    │
        │   ▼                                       │
        │ reconnecting ──重连成功──▶ 恢复中           │
        │                              │            │
        │                    发 backfill(lastSeq)    │
        │                              │            │
        │                    收 backfill-data        │
        │                      ├── hasMore → 继续请求 │
        │                      └── 完成 ──▶ open      │
        └──────────────────────────────────────────┘

4.2 ReconnectingWebSocket 升级(关键增量)

/**
 * 会话恢复管理器:叠在 ReconnectingWebSocket 之上
 * 职责:游标维护、补传请求、补传数据投递
 */
export class SessionManager {
  private lastSeq = 0;                    // 游标:最后收到的快照序号
  private restoring = false;              // 恢复中标记(防补传/实时交错乱序)

  constructor(private conn: ReconnectingWebSocket) {
    // Day 58 预留的钩子今天派上用场
    conn.onSessionRestored = () => this.requestBackfill();

    conn.onMessage((data) => {
      const msg = JSON.parse(data);
      switch (msg.type) {
        case "snapshot":
          if (!this.restoring) {
            // 正常态:直接投递 + 推进游标
            this.lastSeq = msg.seq;
            this.onSnapshot?.(msg.payload);
          } else {
            // 恢复中收到的实时数据:暂存,等补传完成后统一排序投递(见 4.3)
            this.pendingDuringRestore.push(msg);
          }
          break;

        case "backfill-data":
          this.handleBackfill(msg);
          break;
      }
    });
  }

  /** 恢复入口:请求从游标之后的数据 */
  private requestBackfill(): void {
    if (this.lastSeq === 0) return;   // 从未收到过数据(首连)——无需补传
    this.restoring = true;
    this.pendingDuringRestore = [];
    this.conn.send(JSON.stringify({ type: "backfill", fromSeq: this.lastSeq }));
  }

  /** 处理补传批次 */
  private handleBackfill(msg: {
    snapshots: { seq: number; payload: PanelSnapshot }[];
    hasMore: boolean;
    gapLost: boolean;
  }): void {
    if (msg.gapLost) {
      console.warn("部分历史已超出服务器缓冲窗口,无法完整补传");
    }
    // 补传数据按序投递 + 推进游标
    for (const s of msg.snapshots) {
      if (s.seq <= this.lastSeq) continue;   // 防重复(见第五节)
      this.lastSeq = s.seq;
      this.onSnapshot?.(s.payload);
    }
    if (msg.hasMore) {
      // 还有更多:继续翻页请求
      this.conn.send(JSON.stringify({ type: "backfill", fromSeq: this.lastSeq }));
    } else {
      // 补传完成:把恢复期间暂存的实时数据排序补投
      this.pendingDuringRestore
        .sort((a, b) => a.seq - b.seq)
        .forEach((m) => {
          if (m.seq > this.lastSeq) {
            this.lastSeq = m.seq;
            this.onSnapshot?.(m.payload);
          }
        });
      this.restoring = false;
    }
  }

  /** 快照回调(面板层订阅这个) */
  onSnapshot: ((s: PanelSnapshot) => void) | null = null;

  private pendingDuringRestore: { seq: number; payload: PanelSnapshot }[] = [];
}

4.3 为什么恢复期间要暂存实时数据

补传请求在网络上的同时,服务器的新 snapshot 也在推——两条流交错到达。不暂存直接投递,图表数据就会乱序(新点先画、旧点后画,折线打结)。暂存 + 补传完成后统一排序,保证投递顺序 = seq 顺序


五、去重:at-least-once 的业务代价

5.1 哪里会产生重复

1. 补传边界:backfill(fromSeq) 的语义是"严格大于 fromSeq"——若语义含糊(>=),边界重复
2. 恢复期间暂存的实时消息,可能与补传批次尾部重叠
3. (MQTT 场景)QoS 1 重发

5.2 去重的唯一正确姿势:游标比较

// 核心一行:只接受比游标新的数据
if (s.seq <= this.lastSeq) continue;

// ❌ 错误思路:用 Set<seq> 记录"见过的序号"
//    内存无限增长 + 恢复后集合膨胀——去重要利用"单调递增"这个结构性质

单调序列去重 = 维护最大值,O(1) 空间。这是 Day 20 协变逆变那种"利用类型系统结构性质"的思想在网络层的镜像:利用数据的结构性质,算法复杂度断崖下降


六、完整验证实验

今天的验收实验(截图进博客):

1. 启动服务器 + 面板,正常收数(游标推进)
2. DevTools → Network → Offline(或杀服务器进程)
3. 观察面板:状态灯变黄(reconnecting)、曲线停在断线时刻
4. 等待 30 秒(期间服务器数据继续进环形缓冲)
5. 恢复网络(或重启服务器)
6. 验收点:
   a. 状态灯变绿
   b. 曲线上断线缺口被数据填平(无空白、无打结)
   c. 控制台无重复 seq 投递
7. 极限测试:断线 15 分钟(超过缓冲窗口 10 分钟)
   → 面板提示"部分历史不可恢复",曲线留一小段空白但实时流正常继续

七、常见坑点

坑 1:backfill 语义含糊导致边界重复

客户端发 fromSeq=100,服务器回 seq>=100(把 100 又发一遍)。修复:协议文档写死"严格大于",两端实现对照测试。

坑 2:补传大帧卡死连接

缺口 10 分钟 = 600 条一次性发 → 单帧几 MB,弱网下重连都可能被这个帧拖死。修复:分页补传(每批 120 条 + hasMore 翻页,见 3.3/4.2)。

坑 3:恢复期间实时数据抢先投递

不暂存直接画 → 折线打结。修复:restoring 标志 + 暂存队列(4.3)。

坑 4:游标用时间戳

时钟偏差 + 同毫秒多消息 → 边界漏/重。修复:序号游标(2.1)。

坑 5:环形缓冲 shift 的隐藏 O(n)

万级容量的 RingBuffer 每次 push 都 shift → O(n) 摊还。百级容量无所谓;万级改用首尾指针(head/tail index,写满移动 head)。

坑 6:首连也发 backfill

lastSeq=0 时发 backfill(0) → 服务器回全量历史(浪费)。修复:lastSeq 为 0 判定为首连,跳过补传(4.2 第一行)。


八、自测挑战

T1 · 补传全链路(60 分钟)

实现 RingBuffer + 服务器 backfill 处理 + SessionManager,跑通第六节实验。今天核心作业,也是 BOSS 战容错层的直接预演

T2 · 缺口可视化(20 分钟)

面板上把补传回来的数据段画成不同颜色(或半透明),直观验证补传区间正确——调试神器。

T3 · 分页压测(25 分钟)

断线 8 分钟(≈480 条缺口),每批 120 条分页补传。观察:补传期间页面不卡(帧率不掉),补传总耗时记录进博客。

T4 · 乱序注入测试(进阶,30 分钟)

改造服务器:补传响应故意延迟 2 秒发送,验证"恢复期间实时数据暂存 + 最终排序"的正确性——模拟真实网络的乱序到达。


九、总结

子问题 解法 关键点
知道补哪段 序号游标 lastSeq 单调 ID 做游标,时间戳只管显示
历史从哪来 服务器环形缓冲 定长、写满覆盖、分页回放
无缝衔接 restoring 标志 + 暂存排序 补传期间实时流不抢先
重复数据 seq <= lastSeq 则丢弃 单调序列去重 = 维护最大值

至此容错三层(心跳判死 → 退避重连 → 补传衔接)全部完成。明天转向另一个维度:连接健康、数据不丢之后,消息来得太快怎么办——100Hz 的推送撞上 60Hz 的渲染,中间需要一个聪明的调度器。

评论