【WebSocket+MQTT】day62-message-scheduler

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

Day 62 · 高频消息性能 — 批处理、节流与 rAF 对齐:消息风暴的调度艺术

假设链路已经完美:连接稳定、数据不丢。现在服务器把推送频率提到 20Hz(甚至 100Hz 的振动通道)——每秒 100 次 onmessage,每次都去 setOption 更新图表:渲染队列爆炸、主线程被切成碎片、GPU 来不及画。问题从"网络层"转到**“消费节奏”**:数据到得快,不代表要画得快。今天的核心认知:渲染速率应该对齐屏幕刷新率(60fps),而不是对齐数据到达速率——中间需要一个消息调度器。


目录


一、问题复现:消息风暴如何拖垮页面

1.1 实验代码

// ❌ 反模式:每条消息直接更新图表
ws.onmessage = (ev) => {
  const snapshot = decode(ev.data);
  chart.setOption({ series: [{ data: toPoints(snapshot) }] });   // 每秒 100 次!
};

1.2 症状与病理

症状 病理
页面掉到 10fps setOption 每次触发布局重算 + 重绘,100 次/秒远超渲染预算
输入框打字卡顿 主线程被 onmessage + setOption 轮番占用
内存锯齿状暴涨 每次更新分配新数据数组,GC 频繁触发
后台标签页回来时卡死 rAF 被浏览器暂停但消息堆积(见坑 4)

1.3 一笔渲染预算账

60fps → 每帧 16.6ms 预算
每次 setOption(中型图表)≈ 3~8ms
100 次/秒 setOption ≈ 每秒 300~800ms 渲染开销 —— 超预算 5 倍

正确目标:每秒最多 60 次 setOption(且每帧 ≤ 1 次)

二、三种调度策略对比

策略 原理 延迟 问题
❌ 直接消费 每条消息立即处理 0 消息速率 > 渲染速率 → 崩
⚠️ setTimeout 节流 每隔 N ms 处理"最新一条" N ms 与屏幕节拍不同步,可能撕裂/空帧
rAF 对齐 + 批处理 消息先进队列,每次 rAF 取"一批最新"消费 ≤16.6ms 无(今天的方案)

2.1 与旧知识的连接

  • Day 34 的 rAF + dt 帧率无关动画:渲染节拍器思想完全一致
  • Day 56 的分频更新(KPI 5 帧一次、雷达 30 帧一次):不同模块消费频率不同——今天把它泛化成调度器

三、批处理:攒一帧画一次

3.1 核心思想

消息到达:  m1  m2  m3  m4  m5  m6  m7  m8  ...(100/秒)
                │
                ▼ 进队列(不处理)
rAF 节拍:  ──────|─────────|─────────|───(60/秒)
                   │         │
                   ▼         ▼
每帧消费:  [m1..m5 聚合]  [m6..m11 聚合]   ← 攒一批,只画一次

3.2 "聚合"的语义按数据类型分

数据类型 聚合策略 例子
最新值语义 取批内最后一条 当前温度、设备状态、KPI
累计值语义 取批内最后一条(服务器算好累计) 总产量
流式序列 全批追加进滚动窗口 温度曲线(每条都是新的时序点)
事件语义 全批合并处理 告警列表(一条不能丢)

🎯 "取最后"还是"全保留"不是技术选择,是业务语义选择——温度 1 秒内 100 个读数画曲线要全保留(否则曲线失真),但"当前温度显示"只要最新的。


四、rAF 对齐:把消费绑到屏幕节拍

/**
 * rAF 对齐的消息消费循环
 * 关键:队列只负责攒,rAF 负责节拍,消费函数负责取
 */
class RafAlignedConsumer {
  private queue: PanelSnapshot[] = [];
  private rafId = 0;
  private consuming = false;

  constructor(private onBatch: (batch: PanelSnapshot[]) => void) {}

  /** 消息入口:只入队,不做任何处理 */
  push(snap: PanelSnapshot): void {
    this.queue.push(snap);
    this.ensureRunning();
  }

  /** 保证 rAF 循环在跑(队空则停,省电) */
  private ensureRunning(): void {
    if (this.consuming) return;
    this.consuming = true;
    const loop = (): void => {
      if (this.queue.length === 0) {
        this.consuming = false;      // 队空停表:没数据时不空转 rAF
        return;
      }
      const batch = this.queue;
      this.queue = [];
      this.onBatch(batch);           // 一帧一批
      this.rafId = requestAnimationFrame(loop);
    };
    this.rafId = requestAnimationFrame(loop);
  }

  destroy(): void {
    cancelAnimationFrame(this.rafId);
    this.consuming = false;
    this.queue = [];
  }
}

效果:无论消息 10Hz 还是 1000Hz,onBatch 严格 ≤ 60 次/秒——渲染开销封顶在预算内。


五、数据聚合:100Hz 降到 60fps 可用

高频通道(振动 1kHz)即使批处理,一批 16 个点 × 每帧追加也会让曲线窗口瞬间塞满。进图表前先降采样(Day 54 知识的接入层应用):

/**
 * 批内降采样:100Hz → 每帧最多 2 个点(min-max 抽取,保峰值)
 * 与 Day 54 的 lttb 不同:这里在"生产端"降采样(进图表前),
 * lttb 在"渲染端"降采样——两层降采样各司其职
 */
export function minMaxDownsample(
  batch: PanelSnapshot[],
  maxPoints: number
): PanelSnapshot[] {
  if (batch.length <= maxPoints) return batch;
  const step = batch.length / maxPoints;
  const out: PanelSnapshot[] = [];
  for (let i = 0; i < maxPoints; i += 2) {
    // 每个桶取 min 和 max 两个点:保住波形轮廓和峰值
    const bucket = batch.slice(Math.floor(i * step), Math.floor((i + 2) * step));
    if (bucket.length === 0) continue;
    const temps = bucket.map((b) => b.devices[0].temperature);
    const minIdx = temps.indexOf(Math.min(...temps));
    const maxIdx = temps.indexOf(Math.max(...temps));
    out.push(bucket[Math.min(minIdx, maxIdx)]);
    out.push(bucket[Math.max(minIdx, maxIdx)]);
  }
  return out;
}

双层降采样架构(把本周和上周串起来):

1kHz 传感器流
   │ 服务器端:按秒打包成二进制帧(Day 59)
   ▼
网络层(1 帧/秒,帧内 1000 点)
   │ 客户端接入层:min-max 降采样(本节)→ 每秒 60 点
   ▼
MessageScheduler(rAF 对齐批处理)
   │
   ▼
ECharts 渲染层:sampling: 'lttb'(Day 54)→ 屏幕像素级 2 点/px

六、MessageScheduler 完整实现

把批处理 + rAF 对齐 + 分频 + 降采样组装成 BOSS 战直接使用的调度器:

/**
 * 高频消息调度器:BOSS 战数据消费层的核心
 * 职责:
 * 1. 消息入队(onmessage 只做入队,O(1))
 * 2. rAF 节拍批量消费(渲染开销封顶 60 次/秒)
 * 3. 模块分频(KPI 5 帧一次、雷达 30 帧一次——Day 56 策略的调度器化)
 */
export class MessageScheduler {
  private consumer: RafAlignedConsumer;
  private frame = 0;
  private lastKpiFrame = -1;
  private lastRadarFrame = -1;

  constructor(
    private handlers: {
      onStream: (batch: PanelSnapshot[]) => void;      // 曲线:全批(降采样后)
      onKpi: (snap: PanelSnapshot) => void;             // KPI:每 5 帧最新
      onRadar: (snap: PanelSnapshot) => void;           // 雷达:每 30 帧最新
    }
  ) {
    this.consumer = new RafAlignedConsumer((batch) => this.dispatch(batch));
  }

  /** WS onmessage 入口(保持 O(1) 轻量) */
  push(snap: PanelSnapshot): void {
    this.consumer.push(snap);
  }

  /** 每帧分发:按模块频率路由 */
  private dispatch(batch: PanelSnapshot[]): void {
    this.frame++;
    const latest = batch[batch.length - 1];   // 最新值语义的统一取法

    // 曲线:全批降采样后投递(流式语义,不能只取最新)
    this.handlers.onStream(minMaxDownsample(batch, 4));

    // KPI:5 帧一次
    if (this.frame - this.lastKpiFrame >= 5) {
      this.lastKpiFrame = this.frame;
      this.handlers.onKpi(latest);
    }
    // 雷达:30 帧一次
    if (this.frame - this.lastRadarFrame >= 30) {
      this.lastRadarFrame = this.frame;
      this.handlers.onRadar(latest);
    }
  }

  destroy(): void {
    this.consumer.destroy();
  }
}

与前几日模块的组装关系

// Day 63 BOSS 战的组装预览:
const conn = new ReconnectingWebSocket(wsUrl("dev"));   // Day 58
const session = new SessionManager(conn);               // Day 61
const scheduler = new MessageScheduler({                // Day 62(今天)
  onStream: (batch) => lineChart.appendBatch(batch),
  onKpi: (snap) => kpiCard.update(snap.kpi),
  onRadar: (snap) => radarChart.update(snap),
});

session.onSnapshot = (snap) => scheduler.push(snap);    // 补传后的快照也走调度器

七、压测与验收

7.1 实验设计

服务器把推送频率从 1Hz 逐步提到 100Hz,对比"直接消费"和"调度器消费":

推送频率 直接消费 fps 调度器 fps 调度器每秒 setOption 次数
1Hz 60 60 1
10Hz 55 60 10
50Hz 28 60 ~50(仍 ≤ 60)
100Hz 11 60 60(封顶)

验收标准:100Hz 推送下页面保持 58fps+,打字/点击交互无卡顿。

7.2 测量方法

  • fps:Performance 面板录制 20 秒看 Frames 轨道
  • 消费次数:dispatch 里计数,每秒打印一次
  • 交互延迟:页面上放个输入框,高负载下打字对比

八、常见坑点

坑 1:批处理变成了"只处理最后一条"

所有数据都取 latest —— 曲线只剩每帧一个点,波形完全失真。修复:按语义路由(流式全批、状态取最新,见 3.2)。

坑 2:rAF 循环空转

队列空了 rAF 还在跑(每帧空调用)。修复:队空停表(4 的 consuming 标志)。

坑 3:后台标签页消息堆积成山

切走 5 分钟回来,队列里堆了 3 万条消息——rAF 恢复后要"消化"很久。修复:入队时限长(超过 300 条丢最旧的,状态数据只留最新本来就该丢):

push(snap: PanelSnapshot): void {
  this.queue.push(snap);
  if (this.queue.length > 300) {
    this.queue = this.queue.slice(-300);   // 后台堆积防线
  }
  this.ensureRunning();
}

坑 4:visibilitychange 忘了处理

后台时正确姿势:主动退订或降频(Day 57 T4),回前台时清一次队列再继续——配合坑 3 的限长双保险。

坑 5:调度器和补传抢一个入口

Day 61 的补传快照和实时快照走不同路径分别 setOption → 图表更新交错闪烁。修复:统一进调度器(6 的组装代码,session.onSnapshot 是唯一入口)。

坑 6:dispatch 里做重活

分发函数里 JSON.parse 大对象/深拷贝 → 每帧成本失控。纪律:onmessage 只入队,dispatch 只路由,重活全部在各模块自己的 update 里


九、自测挑战

T1 · 调度器 + 压测(50 分钟)

实现 RafAlignedConsumer + MessageScheduler,服务器加"频率"参数(URL 查询串控制 1/10/50/100Hz),跑 7.1 的对比表并填入实测数据。今天核心作业,BOSS 战直接复用

T2 · 语义路由改造(25 分钟)

在 handlers 里增加"告警事件"通道:批内告警一条不能丢(全批投递到告警列表),与流式/状态通道并存。构造 100Hz 数据里混入随机告警,验证零丢失。

T3 · 后台恢复测试(20 分钟)

切后台 3 分钟再回来:验证坑 3 的限长生效(队列 ≤300),回来后 1 秒内恢复正常显示且无长时间"消化"卡顿。

T4 · 调度器性能剖析(进阶,30 分钟)

Performance 录制 100Hz 满负载 30 秒:找出 dispatch 路径上最贵的三个函数,写 200 字分析进博客(哪个模块的 update 最贵?为什么?能不能再分频?)。


十、总结

问题 解法 复用知识
消息速率 > 渲染速率 rAF 对齐批处理 Day 34 的渲染节拍器
批内怎么取 按语义路由(流式全批/状态取最新/事件全保留) Day 56 分频思想
高频点塞爆图表 接入层 min-max 降采样 Day 54 双层降采样架构
后台堆积 入队限长 + visibilitychange
与补传抢入口 统一进调度器 Day 61 的 SessionManager 钩子

本周六天的零件全部备齐:连接管理(57)+ 心跳重连(58)+ 二进制通道(59)+ MQTT 视野(60)+ 补传容错(61)+ 消息调度(62)。明天 BOSS 战:全部组装成 WsDataSource,替换上周面板的 Mock 源,兑现那个说过两次的架构承诺。

评论