Skip to content

WebSocket 可靠性与广播

Phase 06 — Real-time · 知识点 04–07:Heartbeat · Reconnect · Broadcast · Room


1. 学习目标

完成本知识点后,你应该能够:

  • 实现应用层 Heartbeat(Ping/Pong 或自定义心跳消息)
  • 设计客户端 Reconnect 策略(指数退避、快照恢复)
  • 使用 Hub 模式实现 Broadcast 广播推送
  • 按 Room(房间/频道)隔离消息,支持按仓库或区域订阅
  • 处理慢客户端与连接泄漏
  • 为 Three.js 多视图场景提供正确的订阅语义

2. 为什么需要

WebSocket 长连接在移动网络、代理超时、服务端重启时随时可能断开。没有 Heartbeat,你无法区分「空闲」与「死连接」;没有 Reconnect 策略,Three.js 场景会停在旧状态;没有 Broadcast 和 Room,要么全员收到无关数据,要么代码里充斥 if device.Warehouse == client.Warehouse 的散落逻辑。

数字孪生典型场景:同一仓库大屏、调度员控制台、移动端巡检——它们需要同一仓库内广播,但不应收到其他仓库的 AGV 更新。


3. 核心概念

3.1 Heartbeat(心跳)

层级机制说明
协议层WebSocket Ping/Pong 帧gorilla 可设 SetPingHandler
应用层JSON {"type":"ping"} / {"type":"pong"}可携带服务端时间,便于 RTT 统计

服务端策略

  • 每 30s 发送 Ping 或应用层 ping
  • 若 N 次未收到 Pong,主动 Close 连接
  • 配合 SetReadDeadline 检测读超时

3.2 Reconnect(重连)

客户端(Three.js 前端)职责

  1. 监听 onclose / onerror
  2. 指数退避:1s → 2s → 4s → ... → max 30s
  3. 重连成功后请求快照(HTTP 或 WS snapshot 消息),再接收增量
  4. 使用 lastEventIdsince 时间戳补漏(可选)

服务端职责

  • 连接建立时推送全量快照
  • 维护 per-room 最近 N 条事件(可选,用于补漏)
  • 幂等:同一 deviceId + timestamp 客户端可去重

3.3 Broadcast(广播)

         ┌── Client A
Hub ─────├── Client B    所有注册客户端收到同一条消息
         └── Client C

Hub 结构:

go
type Hub struct {
    clients    map[*Client]bool
    broadcast  chan []byte
    register   chan *Client
    unregister chan *Client
}

主循环从 broadcast 读取消息,遍历 clients 写入。写失败则 unregister。

3.4 Room(房间)

Room "warehouse-A" → [Client1, Client2]
Room "warehouse-B" → [Client3]
  • 客户端连接后发送 {"action":"join","room":"warehouse-A"}
  • 推送 AGV 更新时只向该 room 内客户端广播
  • 一个客户端可加入多个 room(如监控多区域)

4. 基础语法

4.1 应用层心跳

go
const (
    writeWait  = 10 * time.Second
    pongWait   = 60 * time.Second
    pingPeriod = (pongWait * 9) / 10
)

func (c *Client) readPump() {
    c.conn.SetReadLimit(512 << 10)
    c.conn.SetReadDeadline(time.Now().Add(pongWait))
    c.conn.SetPongHandler(func(string) error {
        c.conn.SetReadDeadline(time.Now().Add(pongWait))
        return nil
    })
    for {
        _, message, err := c.conn.ReadMessage()
        if err != nil {
            break
        }
        c.hub.handleMessage(c, message)
    }
    c.hub.unregister <- c
}

func (c *Client) writePump() {
    ticker := time.NewTicker(pingPeriod)
    defer ticker.Stop()
    for {
        select {
        case msg, ok := <-c.send:
            c.conn.SetWriteDeadline(time.Now().Add(writeWait))
            if !ok {
                c.conn.WriteMessage(websocket.CloseMessage, []byte{})
                return
            }
            if err := c.conn.WriteMessage(websocket.TextMessage, msg); err != nil {
                return
            }
        case <-ticker.C:
            c.conn.SetWriteDeadline(time.Now().Add(writeWait))
            if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
                return
            }
        }
    }
}

4.2 Room 广播

go
type Room struct {
    id      string
    clients map[*Client]bool
    mu      sync.RWMutex
}

func (r *Room) Broadcast(msg []byte) {
    r.mu.RLock()
    defer r.mu.RUnlock()
    for c := range r.clients {
        select {
        case c.send <- msg:
        default:
            // 发送缓冲区满,视为慢客户端,关闭连接
            close(c.send)
        }
    }
}

4.3 客户端重连伪代码(Three.js)

javascript
let ws;
let retry = 0;
const maxRetryDelay = 30000;

function connect() {
  ws = new WebSocket('wss://api.example.com/ws');
  ws.onopen = () => {
    retry = 0;
    ws.send(JSON.stringify({ action: 'join', room: 'warehouse-A' }));
    fetchSnapshot(); // HTTP 获取全量,保证 Three.js 场景一致
  };
  ws.onclose = () => {
    const delay = Math.min(1000 * 2 ** retry, maxRetryDelay);
    retry++;
    setTimeout(connect, delay);
  };
}

5. 代码解析

SetReadDeadline + PongHandler:每次收到 Pong 刷新 deadline,超时则 ReadMessage 返回错误,触发清理。

c.send 带缓冲 channel:避免慢客户端阻塞 Hub。select default 丢弃或关闭慢连接。

Room 与 Hub 关系:Hub 管理全局注册;Room 是 Hub 内的分组。也可设计为 map[string]*Room

指数退避:避免服务端故障时大量客户端同时重连造成「惊群」。


6. JavaScript / TypeScript 对比

概念前端Go 服务端
心跳定时 ws.send(ping) 或依赖浏览器 PongPingMessage + SetPongHandler
重连onclose + setTimeoutN/A(服务端被动接受新连接)
广播通常不实现,只接收Hub/Room 遍历 Write
房间发送 join 消息维护 room → clients 映射

前端开发者优势:重连 UX 你们更熟悉;本阶段重点是把服务端 Hub/Room 设计正确。


7. 常见错误

错误 1:只在客户端做心跳

代理/NAT 可能在无流量时断开连接,服务端必须主动 Ping。

错误 2:Broadcast 时持锁 Write

在全局锁内写网络 IO 会阻塞所有客户端注册/注销。应只锁数据结构,Write 通过 channel 交给各 client 的 writePump。

错误 3:重连后不拉快照

Three.js 场景只收到增量,AGV 可能「瞬移」或缺失。

错误 4:Room 不校验权限

任意客户端可 join 任意 room,生产环境需 JWT 或 session 校验。

错误 5:unregister 时 double close channel

关闭 c.send 前应检查是否已关闭,或使用 sync.Once


8. 实际应用

数字孪生 — 多仓库 Room 设计

room:warehouse-shanghai  → 上海仓 AGV + 告警
room:warehouse-beijing   → 北京仓
room:alerts-global       → 全局告警(可选)

消息流

  1. AGV 模拟器产生位置 → Pipeline 写入
  2. device.warehouseId 路由到对应 Room
  3. Room.Broadcast → 各客户端 writePump → Three.js 更新 Mesh

Reconnect 流程

断开 → 退避重连 → join room → GET /api/devices/snapshot?warehouse=xxx → 渲染 → 接收增量

9. 深入理解

9.1 惊群与 jitter

重连延迟加随机 jitter(±20%),分散同时重连峰值。

9.2 背压与 drop 策略

实时场景优先最新状态而非完整历史。慢客户端可只保留最后一条位置更新(coalesce)。

9.3 Redis Pub/Sub 扩展

多实例部署时,各实例 Hub 通过 Redis Pub/Sub 同步 Room 消息(Phase 05 已学基础,Project 02 可扩展)。


10. 练习

详细练习见 exercises/phase-06-realtime/02-websocket-reliability-broadcast.md

Level 1 — 基础

练习 1.1:为 WebSocket 服务添加 Ping/Pong,60s 无响应断开。

练习 1.2:实现 join / leave room 消息处理。

Level 2 — 应用

练习 2.1:Hub + Room:向 warehouse-A 广播,其他 room 客户端不应收到。

练习 2.2:写一个简单的 HTML/JS 重连测试页,验证指数退避。

练习 2.3:连接建立时 HTTP 快照 + WS 增量联动。

Level 3 — 综合

练习 3.1:慢客户端模拟:人为延迟 Write,验证 Hub 踢出逻辑。

练习 3.2:支持客户端同时 join 两个 room。

Level 4 — 项目实践

练习 4.1:Project 02 — 完整 Heartbeat + Room Broadcast + 前端 Reconnect 演示。


11. 学习检查

  1. 应用层心跳与协议层 Ping 有何区别?何时两者都用?
  2. 重连后为什么必须拉快照?
  3. Broadcast 时如何避免全局锁阻塞?
  4. Room 与 Redis Channel 命名如何对应?
  5. 什么是惊群,如何缓解?

12. 下一步

已完成下一知识点关系
Heartbeat · Reconnect · Broadcast · RoomConcurrent Connections · Real-time Data Pipeline单连接可靠 → 高并发与数据管道

下一步学习如何支撑上千并发 WebSocket 连接,以及从数据源到前端的完整实时管道架构。


学习导航

上一篇:WebSocket 基础 · 对应练习 · 下一篇:实时数据管道