【WebSocket+MQTT】day63-boss-realtime-panel

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

Day 63 · BOSS 战:实时监控面板 — 把六天的零件拧成一条生产级数据链路

上周的面板跑在定时器模拟数据上,你留了一个 DataSource 接口和一句承诺:“第 9 周换真实数据源,接口不变、零改动”。今天兑现承诺的时刻:本地搭起 WebSocket 服务器模拟设备推送,客户端用本周六天的全部零件(连接管理、心跳重连、二进制通道、补传容错、消息调度)组装 WsDataSource——上周面板代码一行不改,从"模拟"升级为"实时"。最后杀掉服务器,看面板自动重连、断线缺口被补传填平——这是本周学习成果的毕业礼。


目录


一、目标与验收标准

1.1 系统全景

┌─────────────────── 服务器(Node) ───────────────────┐
│  设备模拟器(1Hz 快照 + 偶发告警 + 10Hz 二进制通道)      │
│  环形历史缓冲(600 条)                                  │
│  WS 端点 /sensor-feed(backfill 分页补传)               │
└──────────────────────┬───────────────────────────────┘
                       │ WebSocket(心跳 10s/超时 5s)
┌──────────────────────▼───────────────────────────────┐
│  ReconnectingWebSocket(指数退避重连)                   │
│  SessionManager(游标 + 补传 + 暂存排序)                │
│  MessageScheduler(rAF 批处理 + 分频 + 降采样)          │
│  ── PanelSnapshot(上周的数据契约,原样复用)──          │
└──────────────────────┬───────────────────────────────┘
                       │ DataSource 接口
┌──────────────────────▼───────────────────────────────┐
│  上周的 PanelApp:五模块面板 —— 一行不改                 │
└──────────────────────────────────────────────────────┘

1.2 硬性验收指标

指标 标准 检验方法
接口零改动 PanelApp 及五个图表模块 diff 为零 git diff 验证
实时性 快照 1s 一推,曲线平滑滚动 肉眼 10 分钟
断线恢复 杀服务器 → 状态灯黄 → 重启 → 缺口填平无缝衔接 第六节实验
帧率 10Hz 推送下 ≥ 58fps Performance 录制
内存 30 分钟 JS Heap 稳定 堆快照对比
生命周期 页面关闭后服务器连接数归零 服务器日志

二、系统架构总览

2.1 目录结构(本周工程的最终形态)

realtime-panel/
├── server/                    # 服务器端(仅开发环境运行)
│   ├── device-simulator.ts    # 设备模拟器(快照 + 二进制振动通道)
│   ├── ring-buffer.ts         # Day 61 环形缓冲
│   └── server.ts              # WS 端点 + backfill 处理
├── src/                       # 浏览器端
│   ├── types.ts               # PanelSnapshot / DataSource 契约(上周原文件)
│   ├── ws/
│   │   ├── reconnecting-ws.ts # Day 58
│   │   ├── heartbeat.ts       # Day 58
│   │   ├── session.ts         # Day 61 SessionManager
│   │   └── scheduler.ts       # Day 62 MessageScheduler
│   ├── binary/
│   │   └── codec.ts           # Day 59 encode/decode(双端共用)
│   ├── data/
│   │   ├── mock-source.ts     # 上周 MockDataSource(保留:单测用)
│   │   └── ws-source.ts       # ⭐ 今天的主角
│   ├── charts/                # 上周五模块(一行不改)
│   ├── panel-app.ts           # 上周 PanelApp(一行不改)
│   └── main.ts                # 入口:只换 DataSource 一行
└── week09 计划文档...

2.2 数据流分工(本周六天的落位)

模块 职责
传输 ReconnectingWebSocket 57/58 连接、心跳判死、退避重连
会话 SessionManager 61 游标、补传、乱序防护
消费 MessageScheduler 62 rAF 批处理、分频、降采样
数据源 WsDataSource 今天 把以上组装成 DataSource 契约
视图 PanelApp + charts 上周 五模块面板(零改动)

三、服务器端:设备模拟器

3.1 完整服务器(Day 57/59/61 服务器能力的总装)

/**
 * 实时面板服务器:快照流 + 二进制振动通道 + 历史补传
 * 运行:npx tsx server/server.ts
 */
import { WebSocketServer, WebSocket } from "ws";
import { RingBuffer } from "./ring-buffer";
import { encodeChannelFrame } from "../src/binary/codec";
import type { PanelSnapshot } from "../src/types";

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

console.log("实时面板服务器已启动 ws://localhost:8080/sensor-feed");

wss.on("connection", (ws: WebSocket) => {
  console.log("客户端接入,当前连接数:", wss.clients.size);

  // ── 1Hz 快照流(JSON:低频控制流用 JSON,Day 59 的选型原则)──
  const snapTimer = setInterval(() => {
    if (ws.readyState !== WebSocket.OPEN) return;
    tick++;
    const payload = generateSnapshot(tick);
    const item = { seq: ++seq, payload };
    history.push(item);
    ws.send(JSON.stringify({ type: "snapshot", seq: item.seq, payload }));
  }, 1000);

  // ── 10Hz 二进制振动通道(高频通道用二进制,每 10 帧打一包)──
  const vibTimer = setInterval(() => {
    if (ws.readyState !== WebSocket.OPEN) return;
    // 每包 10 个 float32 采样(10Hz 聚合发送,模拟真实网关的打包行为)
    const values = new Float32Array(10);
    for (let i = 0; i < 10; i++) {
      values[i] = Math.sin((tick * 10 + i) / 8) * 2 + (Math.random() - 0.5);
    }
    ws.send(encodeChannelFrame(1, Date.now(), 100, 0x01, values));
  }, 1000);

  // ── 消息处理(心跳 + 补传)──
  ws.on("message", (raw) => {
    const msg = JSON.parse(raw.toString());
    if (msg.type === "ping") {
      ws.send(JSON.stringify({ type: "pong", seq: msg.seq }));
    } else if (msg.type === "backfill") {
      const { items, hasMore, gapLost } = history.queryAfter(msg.fromSeq, 120);
      ws.send(JSON.stringify({
        type: "backfill-data",
        fromSeq: msg.fromSeq,
        snapshots: items,
        hasMore,
        gapLost,
      }));
    }
  });

  ws.on("close", () => {
    clearInterval(snapTimer);
    clearInterval(vibTimer);
    console.log("客户端断开,剩余连接数:", wss.clients.size);
  });
});

/**
 * 快照生成(与上周 MockDataSource 同款仿真模型——保证面板行为一致)
 */
function generateSnapshot(t: number): PanelSnapshot {
  const temp = (base: number) => {
    const wave = 5 * Math.sin(t / 40 + base);
    const noise = (Math.random() - 0.5) * 2;
    const spike = Math.random() < 0.0005 ? 12 : 0;   // 万分之五概率尖峰触发告警色
    return Number((base + wave + noise + spike).toFixed(1));
  };
  return {
    ts: Date.now(),
    devices: [
      { deviceId: "f1", name: "1号炉", temperature: temp(80), pressure: 2.1 + Math.random() * 0.2, health: [92, 68, 85, 74, 90] },
      { deviceId: "f2", name: "2号炉", temperature: temp(78), pressure: 2.0 + Math.random() * 0.2, health: [88, 75, 80, 82, 76] },
      { deviceId: "f3", name: "3号机组", temperature: temp(82), pressure: 2.3 + Math.random() * 0.2, health: [90, 72, 88, 79, 92] },
    ],
    kpi: {
      totalOutput: 125800 + t * 3,
      alarmCount: Math.random() < 0.03 ? 1 : 0,
      oee: 76 + Math.sin(t / 60) * 4,
      onlineRate: 98 + Math.random() * 1.5,
    },
    energyRanking: [
      { name: "3车间", value: 98 + Math.random() * 3 },
      { name: "1车间", value: 85 + Math.random() * 3 },
      { name: "4车间", value: 72 + Math.random() * 3 },
      { name: "2车间", value: 60 + Math.random() * 3 },
    ],
  };
}

💡 快照模型故意与上周 Mock 一致:这样面板上线后行为完全相同,唯一变化是数据来源——这是验证"接口抽象正确性"的最佳方式:换个实现,视图层毫无感知。


四、WsDataSource:兑现承诺的一刻

4.1 核心类

import type { DataSource, PanelSnapshot } from "../types";
import { ReconnectingWebSocket } from "../ws/reconnecting-ws";
import { SessionManager } from "../ws/session";
import { MessageScheduler } from "../ws/scheduler";
import { decodeChannelFrame } from "../binary/codec";
import { wsUrl } from "../ws/url";

/**
 * WebSocket 数据源:实现上周定义的 DataSource 契约
 * 组装链:ReconnectingWebSocket → SessionManager → MessageScheduler
 *
 * 面板视角:它与 MockDataSource 唯一的区别是实现类名。
 */
export class WsDataSource implements DataSource {
  private conn: ReconnectingWebSocket;
  private session: SessionManager;
  private scheduler: MessageScheduler;
  private listeners: ((s: PanelSnapshot) => void)[] = [];

  constructor(
    private opts = {
      url: wsUrl("dev"),
      // 二进制高频通道的消费者(面板可选实现;不传则丢弃)
    }
  ) {
    // ── 传输层:心跳 + 重连 ──
    this.conn = new ReconnectingWebSocket(this.opts.url);

    // ── 会话层:补传 + 乱序防护 ──
    this.session = new SessionManager(this.conn);

    // ── 消费层:rAF 批处理 + 分频 ──
    this.scheduler = new MessageScheduler({
      onStream: (batch) => this.emitBatch(batch),      // 曲线:全批(降采样后)
      onKpi: (snap) => this.emit(snap),                // KPI:每 5 帧最新
      onRadar: (snap) => this.emit(snap),              // 雷达:每 30 帧最新
    });

    // 快照统一入口:补传恢复的 + 实时的,都走 SessionManager → 调度器
    this.session.onSnapshot = (snap) => this.scheduler.push(snap);

    // 二进制振动帧(独立通道,可选消费)
    this.conn.onBinary?.((buf: ArrayBuffer) => {
      const frame = decodeChannelFrame(buf);
      if (frame) this.onChannelFrame?.(frame);
    });
  }

  /** DataSource 契约:订阅 */
  subscribe(listener: (s: PanelSnapshot) => void): () => void {
    this.listeners.push(listener);
    return () => { this.listeners = this.listeners.filter((l) => l !== listener); };
  }

  /** DataSource 契约:启动 */
  start(): void {
    this.conn.connect();
  }

  /** DataSource 契约:停止(干净关闭,不触发重连) */
  stop(): void {
    this.scheduler.destroy();
    this.conn.close();
  }

  /** 连接状态(面板状态灯订阅——Day 58 T3 的组件) */
  onState(listener: (s: "connecting" | "open" | "reconnecting" | "closed") => void): () => void {
    return this.conn.onState(listener);
  }

  /** 二进制帧回调(振动通道扩展图表用) */
  onChannelFrame: ((f: { channelId: number; values: Float32Array; startTs: number }) => void) | null = null;

  // ── 内部分发 ──
  private emit(snap: PanelSnapshot): void {
    this.listeners.forEach((l) => l(snap));
  }
  private emitBatch(batch: PanelSnapshot[]): void {
    for (const s of batch) this.emit(s);
  }
}

4.2 main.ts:改动的全部代码

// 上周的 main.ts,改动只有这一行:
- const source = new MockDataSource();
+ const source = new WsDataSource();     // ⭐ 整个"模拟→实时"升级的全部 diff

git diff 给自己看:视图层零改动——上周埋的接口设计今天结账。这就是第 1 周 ES Module 分层、第 8 周 DataSource 契约的复利。


五、组装与运行

5.1 启动流程

# 终端 1:服务器
npx tsx server/server.ts

# 终端 2:前端(Vite 开发服务器)
npm run dev

5.2 首次验收检查点

1. 面板加载 → 状态灯绿 → 1 秒后曲线开始滚动 ✅
2. DevTools → Network → WS:
   a. Headers:握手成功(101)
   b. Messages:每秒 1 条 snapshot + 每秒 1 条二进制帧(Binary 标记)
   c. 每 10 秒一对 ping/pong
3. KPI 翻牌器每 5 秒滚动、雷达约 30 秒变一次(分频生效)
4. 控制台无报错、无重复 seq

六、容错验收:亲手杀服务器

6.1 三幕剧(全程录屏进博客)

第一幕:谋杀

1. 面板正常跑 1 分钟(游标推进到 ~60)
2. 服务器终端 Ctrl+C(或任务管理器杀进程)

第二幕:黑暗中重生

3. 客户端观察:
   - 已发的帧收完 → 心跳 ping 发出 → 5 秒无 pong → 判死
   - 状态灯变黄(reconnecting),控制台打印退避序列
   - 面板曲线冻结在最后时刻(不清空、不报错——降级体验)
4. 等待期间服务器历史缓冲虽然没了(进程死了),
   但重启后重新灌数据——这里演示的是"重启后补传窗口从零开始",
   真实生产中历史缓冲在独立进程/数据库里不受影响

第三幕:复活与填坑

5. 重启服务器
6. 客户端退避周期到点 → 重连成功 → session.onSessionRestored 触发
   → backfill(lastSeq) → 服务器从游标回放(新进程 seq 从 1 开始
   → lastSeq=60 > 新 seq → 服务器返回空 + 无更多 → 恢复完成)
   ⚠️ 注意:本实验因服务器重启 seq 重置,补传拿不到旧数据——
   这暴露了"seq 必须持久化"的生产课题(见 T2)
7. 曲线从重连时刻无缝继续(时间轴无跳变)

6.2 不重启的断网版(补传的完整演示)

1. DevTools → Network → Offline(服务器进程不死!)
2. 等 30 秒(服务器缓冲持续积累 30 条)
3. 恢复在线
4. 验收:重连 → backfill 拿回 30 条 → 曲线缺口被填平,
   补传段用 T2 的可视化染色 → 一段颜色不同的补丁清晰可见 ✅

七、性能验收

场景 指标 目标 实测(参考)
稳态(1Hz 快照 + 10Hz 振动帧) fps ≥58 60
稳态 每秒 setOption 次数 ≤60 ~12(调度器分频生效)
稳态 30 分钟 JS Heap 稳定 ±2MB
关页后 服务器连接数 0 0(close 日志确认)

八、常见坑点

坑 1:服务器重启后 seq 归零,游标判断失效

backfill(fromSeq=60) 对上新进程的 seq 1~30 → 服务器查不到 > 60 的数据 → 返回空。行为上"看起来正常"(恢复到实时流),但缺口数据静默丢失。生产修复:seq 用持久化序列(Redis INCR / 数据库自增)。BOSS 战范围内:接受丢失 + 控制台 warn 提示(Day 61 的 gapLost 语义)。

坑 2:振动二进制帧和快照 JSON 抢 onmessage

typeof ev.data 分流漏了分支 → 二进制帧进 JSON.parse 报错刷屏。修复:Day 57 的分流纪律 + onerror 不 print 栈。

坑 3:组件热更新(Vite HMR)后双连接

HMR 重新执行 main.ts,旧 DataSource 没 stop → 服务器日志出现两条连接、面板数据翻倍。修复:HMR 钩子里 import.meta.hot.dispose(() => source.stop())

坑 4:stop() 顺序错误

conn.close()scheduler.destroy() → close 后残余消息还在 rAF 队列里被消费(对已卸载图表 setOption 报错)。修复:先停消费(scheduler),再停会话,最后关连接——从下游往上游关。

坑 5:状态灯不订阅就永远不亮

WsDataSource 的 onState 要在 start() 之前订阅(否则错过 connecting 状态)。习惯:所有事件订阅都在启动前完成


九、自测挑战

T1 · 全链路跑通(90 分钟)

按第三~五节完成服务器 + WsDataSource + 面板接入,git diff 证明视图层零改动,跑通 5.2 检查点。今天核心作业

T2 · 补传染色可视化(20 分钟)

曲线模块加"补传数据段半透明"标记(快照里加个 backfilled: boolean 字段)——第六节实验的杀手级演示截图。

T3 · seq 持久化(进阶,40 分钟)

服务器 seq 写本地文件(每次自增后 fs.writeFile),重启后从持久化值继续。验证:重启实验中缺口被真正补回(观察染色段)。

T4 · MQTT 双源实验(进阶,45 分钟)

按 Day 60 T4 的 MqttDataSource 也接入面板(加个开关切换 WS/MQTT 源)。用 MQTTX 手动发 Topic 数据驱动面板——验证 DataSource 抽象对第二种真实协议依然成立,这是架构含金量的终极测试。

T5 · 弱网模拟(30 分钟)

DevTools → Network → Slow 3G + 丢包(Custom 规则):观察心跳超时误杀率、补传分页耗时。写下:弱网下应该调整哪三个参数?(提示:心跳间隔、退避上限、分页大小)


十、总结与第 9 周 / 第 2 个月毕业

10.1 本周零件的最终落位

零件 出自 在 WsDataSource 里的角色
ReconnectingWebSocket Day 57/58 传输层:连接 + 心跳 + 重连
codec(encode/decode) Day 59 二进制振动通道
SessionManager Day 61 会话层:游标 + 补传
MessageScheduler Day 62 消费层:批处理 + 分频
PanelApp + 五图表 第 8 周 视图层:零改动
DataSource 契约 第 8 周 Day 56 全部的粘合剂

10.2 第 2 个月(Canvas 实时可视化)毕业检查

  • [ ] 流程图编辑器(Week 7 BOSS 战)
  • [ ] 图表面板 + 10 万点性能报告(Week 8)
  • [ ] 实时监控面板 + 断线补传录屏(本周)
  • [ ] 三份作品的 GitHub 仓库 + README + 博客
  • [ ] 能白板讲清:一次温度数据从传感器到像素的完整旅程(MQTT→网关→WS→调度器→ECharts)

10.3 下周预告:第 10 周(Day 64-70)月度综合项目

第 2 个月的收官周——智慧工厂大屏雏形,把两个月的全部能力整合成一个可展示作品:

  • Day 64-65:项目架构与数据层整合(上周面板升级为多屏路由)
  • Day 66-67:大屏视觉工程(标题栏/边框组件/动效细节/自适应复查)
  • Day 68-69:整合 Canvas 编辑器(拓扑图页)+ 实时面板(监控页)
  • Day 70:BOSS 战——完整大屏 Demo + 演示视频 + 第 2 个月毕业总结

这个 Demo 就是求职简历上"项目经验"的第一条,也是阶段 2(0-3 个月智慧工厂大屏)的正式开工。


第 9 周毕业,第 2 个月毕业。从 Day 29 的第一个 canvas 像素,到今天这条心跳、重连、补传俱全的实时链路——你已经具备工业可视化最核心的两条腿:会画(Canvas/图表),会连(实时通信)。下周把它们拧成一件拿得出手的作品。

评论