外观
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 前端)职责:
- 监听
onclose/onerror - 指数退避:
1s → 2s → 4s → ... → max 30s - 重连成功后请求快照(HTTP 或 WS
snapshot消息),再接收增量 - 使用
lastEventId或since时间戳补漏(可选)
服务端职责:
- 连接建立时推送全量快照
- 维护 per-room 最近 N 条事件(可选,用于补漏)
- 幂等:同一
deviceId + timestamp客户端可去重
3.3 Broadcast(广播)
┌── Client A
Hub ─────├── Client B 所有注册客户端收到同一条消息
└── Client CHub 结构:
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) 或依赖浏览器 Pong | PingMessage + SetPongHandler |
| 重连 | onclose + setTimeout | N/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 → 全局告警(可选)消息流:
- AGV 模拟器产生位置 → Pipeline 写入
- 按
device.warehouseId路由到对应 Room - 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. 学习检查
- 应用层心跳与协议层 Ping 有何区别?何时两者都用?
- 重连后为什么必须拉快照?
- Broadcast 时如何避免全局锁阻塞?
- Room 与 Redis Channel 命名如何对应?
- 什么是惊群,如何缓解?
12. 下一步
| 已完成 | 下一知识点 | 关系 |
|---|---|---|
| Heartbeat · Reconnect · Broadcast · Room | Concurrent Connections · Real-time Data Pipeline | 单连接可靠 → 高并发与数据管道 |
下一步学习如何支撑上千并发 WebSocket 连接,以及从数据源到前端的完整实时管道架构。