WebSocket 的测试难点不在于"能连上、能收发",而在于那些只在时间维度上才暴露的问题。 一次请求-响应的 REST 测试是"瞬时快照",而 WebSocket 是一条持续数小时的长连接——它会在网络切换时断开、会在消息乱序时错乱、会在重连风暴中把服务端打挂、会在客户端慢消费时把内存撑爆。这些问题的共同点是:它们不在单条消息里,而在消息之间的时序、连接的生命周期、以及并发连接的规模里。 本文要回答的是:如何围绕连接生命周期、消息语义、重连恢复、并发规模这四层,构建一套能抓住"实时"特有缺陷的测试体系。
实时系统还有一个容易被忽视的性质:它天生是有状态、有并发的。同一条测试在本地跑十次都通过,在 CI 的慢机器上可能因为时序抖动而偶发失败。因此实时通信测试的纪律,一半在"测什么",另一半在"如何让测试本身足够确定"。
一、实时通信测试的独特性
1.1 与请求-响应测试的本质差异
| 维度 | HTTP 请求-响应 | WebSocket 长连接 | 测试影响 |
|---|---|---|---|
| 状态 | 无状态(每请求独立) | 有状态(连接即会话) | 需测连接级状态机 |
| 时序 | 无(一问一答) | 强时序(消息顺序) | 需测顺序、乱序、丢失 |
| 生命周期 | 请求结束即释放 | 长连接(可能数小时) | 需测心跳、重连、泄漏 |
| 并发 | 无连接上限概念 | 连接数受 fd/内存限制 | 需测连接规模与广播 |
| 失败模式 | 超时/错误码 | 静默断开、半开连接 | 需测"假连接"检测 |
1.2 测试层次
层次一:协议层 —— 握手、升级、子协议协商、帧格式
层次二:连接层 —— 鉴权、心跳、优雅关闭、半开检测
层次三:消息层 —— 顺序、去重、补发、背压
层次四:恢复层 —— 重连退避、状态恢复、重连风暴
层次五:规模层 —— 连接数、广播扇出、内存增长
二、连接生命周期测试
2.1 握手与升级
WebSocket 由 HTTP 升级而来,握手阶段就能拦下大量非法连接:
def test_handshake_requires_upgrade_header(ws_url):
# 缺少 Upgrade 头的普通 GET 不应被升级
resp = requests.get(ws_url, headers={"Connection": "close"})
assert resp.status_code in (400, 426), "应拒绝非升级请求"
def test_subprotocol_negotiation(ws_url):
import websocket
ws = websocket.create_connection(ws_url, subprotocols=["graphql-ws"])
assert ws.getsubprotocol() == "graphql-ws"
ws.close()
2.2 连接级鉴权
鉴权必须在握手阶段完成,且失败要明确关闭连接(而不是升级后再拒绝):
def test_connection_rejected_without_token(ws_url):
import websocket
with pytest.raises(websocket.WebSocketBadStatusException) as exc:
websocket.create_connection(ws_url) # 不带 token
assert exc.value.status_code == 401
def test_connection_rejected_with_expired_token(ws_url, expired_token):
import websocket
with pytest.raises(websocket.WebSocketBadStatusException) as exc:
websocket.create_connection(f"{ws_url}?token={expired_token}")
assert exc.value.status_code == 401
⚠️ 不要把鉴权放到连接建立之后的"第一条消息"里。若握手成功、之后才发消息拒绝,攻击者可以批量建立空连接耗尽服务端资源。测试要断言"未授权连接在握手阶段即被 401 拒绝"。
2.3 心跳与半开连接检测
TCP 长连接最阴险的故障是"半开连接"——一端以为连着,另一端早已断开,但操作系统不会立即报错。心跳(ping/pong)是唯一可靠的检测手段:
// 服务端心跳配置(ws 库)
const wss = new WebSocketServer({ server, clientTracking: true });
const interval = setInterval(() => {
wss.clients.forEach((ws) => {
if (ws.isAlive === false) return ws.terminate(); // 上一轮没回 pong
ws.isAlive = false;
ws.ping(); // 发 ping,等待 pong
});
}, 30000);
wss.on("connection", (ws) => {
ws.isAlive = true;
ws.on("pong", () => { ws.isAlive = true; });
});
对应测试要验证"不回 pong 的连接会被终止":
def test_dead_connection_is_terminated(ws_url):
import websocket, time
ws = websocket.create_connection(ws_url)
# 劫持底层 socket,吞掉 pong 响应(模拟半开)
ws.sock.settimeout(0.01)
# 等超过心跳周期 * 2
time.sleep(70)
# 连接应已被服务端终止
with pytest.raises(Exception):
ws.send("probe")
2.4 优雅关闭
服务端重启或下线时,应先发关闭帧(close frame)再断开,让客户端有机会保存状态:
def test_graceful_shutdown_sends_close_frame(ws_url, server_control):
import websocket
ws = websocket.create_connection(ws_url)
server_control.initiate_shutdown()
# 客户端应收到 close frame,code 1001(going away)
frame = ws.recv_frame()
assert frame.opcode == websocket.ABNF.OPCODE_CLOSE
assert frame.data == b"\x03\xe9" # 1001
三、消息时序与乱序
3.1 顺序保证
如果协议承诺"同一连接内的消息有序",测试必须验证这一点:
def test_message_order_preserved(ws_url):
import websocket, json
ws = websocket.create_connection(ws_url)
ws.send(json.dumps({"type": "subscribe", "channel": "trades"}))
seqs = []
for _ in range(100):
msg = json.loads(ws.recv())
seqs.append(msg["seq"])
assert seqs == sorted(seqs), "消息顺序被破坏"
assert seqs == list(range(1, 101)), "存在消息丢失"
3.2 至少一次语义:重复与丢失
实时系统通常提供 at-least-once(至少一次)投递,意味着消息可能重复、也可能在重连后补发。测试要覆盖"重复消息被幂等处理":
def test_duplicate_messages_are_idempotent(client):
# 模拟服务端重发同一 seq 的消息
client.inject({"seq": 5, "type": "trade", "price": 100})
client.inject({"seq": 5, "type": "trade", "price": 100}) # 重复
assert client.state.trade_count == 1, "重复消息被重复处理"
def test_reconnect_replays_missed_messages(ws_url):
import websocket, json
ws = websocket.create_connection(ws_url)
ws.send(json.dumps({"type": "subscribe", "channel": "trades"}))
first = json.loads(ws.recv())
last_seq = first["seq"]
ws.close()
# 重连时携带 last_seq,服务端应补发遗漏的消息
ws2 = websocket.create_connection(f"{ws_url}?resume_from={last_seq}")
resumed = json.loads(ws2.recv())
assert resumed["seq"] == last_seq + 1, "重连未补发遗漏消息"
3.3 背压与缓冲
当消费端比生产端慢时,消息会在缓冲区堆积。测试要验证背压策略(丢弃/降采样/断连)确实生效:
def test_backpressure_drops_or_disconnects(ws_url):
import websocket
ws = websocket.create_connection(ws_url)
ws.send(json.dumps({"type": "subscribe", "channel": "high_freq"}))
# 故意不 recv,让服务端缓冲堆积
import time
time.sleep(10)
# 服务端应在缓冲超过阈值后采取行动(关闭连接或丢弃旧消息)
ws.sock.settimeout(1)
got = ws.recv() # 要么收到"背压丢弃"通知,要么连接已断
assert got is not None
⚠️ 无界缓冲区是内存泄漏的前身。生产环境必须有界,且策略要显式:超阈值时是"丢弃最旧"(适合行情)还是"断开慢消费者"(适合交易确认)?测试要把这个策略固化成用例,否则它会在一次流量高峰里悄悄变成 OOM。
四、断线重连测试
4.1 指数退避
客户端重连必须使用指数退避(exponential backoff),否则服务端一抖动就会遭遇全网重连:
// 客户端重连策略
function reconnect(attempt) {
const base = 500; // 500ms
const max = 30000; // 上限 30s
const jitter = Math.random() * 0.3; // 30% 抖动
const delay = Math.min(max, base * 2 ** attempt) * (1 + jitter);
setTimeout(connect, delay);
}
测试验证退避序列符合预期:
def test_reconnect_uses_exponential_backoff(client, clock):
delays = []
client.on_reconnect_scheduled(lambda d: delays.append(d))
client.simulate_disconnect()
clock.advance(60000) # 推进虚拟时钟 60s
# 退避应递增且有上限
assert delays[0] < delays[1] < delays[2]
assert all(d <= 30000 * 1.3 for d in delays), "退避超过上限+抖动"
4.2 重连风暴防护
当服务端重启,成千上万客户端会同时重连。测试要模拟这个场景并验证服务端不被冲垮:
def test_reconnect_storm_does_not_crash_server(ws_url):
import websocket, threading
conns = []
# 建立 500 个连接
for _ in range(500):
conns.append(websocket.create_connection(ws_url))
# 全部同时断开
for c in conns:
c.close()
# 全部立即重连
ok = 0
for _ in range(500):
try:
websocket.create_connection(ws_url, timeout=5)
ok += 1
except Exception:
pass
assert ok >= 450, f"重连风暴下成功率过低: {ok}/500"
ℹ️ 重连风暴的工程解法是"抖动 + 服务端限流":客户端退避带随机抖动打散重连时刻,服务端用令牌桶限制单位时间的握手数。测试要同时验证这两侧——只测客户端退避、不测服务端限流,风暴依然会打挂服务。
4.3 状态恢复
重连后,客户端需要恢复订阅与游标。测试验证"重连后状态与断开前一致":
def test_state_restored_after_reconnect(client):
client.connect()
client.subscribe("trades")
client.subscribe("orders")
client.receive_until(seq=10)
snapshot_before = client.subscription_snapshot()
client.simulate_disconnect()
client.reconnect()
assert client.subscription_snapshot() == snapshot_before, "订阅未恢复"
# 恢复后应从断点续传,而非从头
assert client.last_seq == 10
五、负载与并发测试
5.1 连接数规模
WebSocket 的资源瓶颈通常是文件描述符和内存,而非 CPU。压测要找到"单机能维持多少连接":
// k6 WebSocket 连接压测
import ws from "k6/ws";
import { check } from "k6";
export const options = {
stages: [
{ duration: "1m", target: 5000 }, // 1 分钟爬到 5000 连接
{ duration: "3m", target: 5000 }, // 保持
{ duration: "1m", target: 0 },
],
};
export default function () {
const url = "wss://api.example.com/ws";
const res = ws.connect(url, {}, (socket) => {
socket.on("open", () => socket.send(JSON.stringify({ type: "subscribe" })));
socket.on("message", (msg) => { /* 消费消息 */ });
socket.setTimeout(() => socket.close(), 10000);
});
check(res, { "status 101": (r) => r && r.status === 101 });
}
# 观察服务端 fd 与内存
watch -n1 'ls /proc/$(pgrep -f ws-server)/fd | wc -l'
# 连接数应接近目标,且内存随连接数线性增长后趋于平稳
5.2 广播扇出
广播场景的复杂度是"连接数 × 消息频率"。测一个 1 万连接的频道广播:
def test_broadcast_fanout_latency(ws_url):
import websocket, time
subscribers = [websocket.create_connection(ws_url) for _ in range(1000)]
for s in subscribers:
s.send(json.dumps({"type": "subscribe", "channel": "trades"}))
start = time.time()
publisher = websocket.create_connection(ws_url)
publisher.send(json.dumps({"type": "publish", "channel": "trades", "seq": 1}))
received = 0
deadline = start + 5
while time.time() < deadline and received < 1000:
for s in subscribers:
s.sock.settimeout(0.01)
try:
s.recv(); received += 1
except Exception:
pass
assert received == 1000, f"扇出丢失: {received}/1000"
5.3 连接泄漏检测
长跑测试要验证"连接建立-关闭"不会泄漏 fd 或 goroutine/线程:
# 用 /proc 观察 fd 数在反复连接-断开后是否回落
for i in $(seq 1 10); do
# 建立并关闭 1000 个连接
k6 run --vus 1000 --iterations 1000 connect_churn.js
echo "after round $i: $(ls /proc/$(pgrep -f ws-server)/fd | wc -l) fds"
done
# fd 数应在每轮后回落到基线,若持续增长即泄漏
六、Mock 服务端与测试替身
6.1 为什么需要 Mock WebSocket 服务端
前端/客户端的重连、时序逻辑测试,不应该依赖真实后端。一个可控的 mock 服务端可以精确制造"乱序"“延迟"“中途断开”:
// 用 ws 库搭一个可控 mock 服务端
import { WebSocketServer } from "ws";
export function createMockServer(port) {
const wss = new WebSocketServer({ port });
const control = { dropNext: false, delayMs: 0, reorder: false };
wss.on("connection", (ws) => {
ws.on("message", (data) => {
if (control.dropNext) { control.dropNext = false; return; } // 制造丢包
const deliver = () => ws.send(JSON.stringify({ echo: data.toString() }));
if (control.delayMs) setTimeout(deliver, control.delayMs);
else deliver();
});
});
return { wss, control, close: () => wss.close() };
}
test("客户端在消息延迟时保持顺序", async () => {
const { control, close } = createMockServer(9999);
control.delayMs = 50;
const client = new RealtimeClient("ws://localhost:9999");
await client.connect();
client.send("a"); client.send("b");
const received = await client.collect(2);
expect(received).toEqual(["a", "b"]);
close();
});
6.2 录制回放
把生产环境的一次真实会话录制成帧序列,测试时回放,能复现线上才出现的时序:
# 录制:抓取真实会话的帧
def record_session(ws_url, output):
import websocket, json
ws = websocket.create_connection(ws_url)
frames = []
for _ in range(500):
frames.append({"dir": "in", "data": ws.recv(), "t": time.time()})
json.dump(frames, open(output, "w"))
# 回放:按录制的时间戳重放
def replay_session(record_file, client):
frames = json.load(open(record_file))
base = frames[0]["t"]
for f in frames:
time.sleep(max(0, f["t"] - base)) # 或注入虚拟时钟
client.inject(f["data"])
七、CI 集成与常见陷阱
7.1 让实时测试在 CI 里稳定
实时测试天生容易 flaky,必须用确定性的时间与同步原语:
# 反模式:用真实 sleep 等待
time.sleep(2)
assert client.received_count == 10 # 慢机器上会失败
# 正确:等待条件而非固定时间
def wait_until(predicate, timeout=5):
deadline = time.time() + timeout
while time.time() < deadline:
if predicate():
return True
time.sleep(0.01)
raise AssertionError("条件未在超时前满足")
wait_until(lambda: client.received_count == 10)
7.2 常见陷阱对照表
| 陷阱 | 现象 | 对策 |
|---|---|---|
| 鉴权放到连接后 | 空连接耗尽资源 | 握手阶段 401 |
| 无心跳检测 | 半开连接长期占用 | ping/pong + 超时终止 |
| 重连无退避 | 服务端抖动引发风暴 | 指数退避 + 抖动 |
| 重连风暴无防护 | 服务端重启后被打挂 | 客户端抖动 + 服务端限流 |
| 无界缓冲 | 慢消费者导致 OOM | 有界缓冲 + 显式策略 |
| 用 sleep 等待 | CI 上偶发失败 | 条件等待 + 虚拟时钟 |
| 只测单连接 | 扇出/规模问题漏测 | 连接数 + 广播压测 |
八、总结
实时通信与 WebSocket 测试的独特性,在于它把风险从"单条消息的内容"转移到了"连接的生命周期、消息之间的时序、以及并发连接的规模”。测试要覆盖握手鉴权与心跳保活(连接层)、顺序与至少一次语义(消息层)、退避与状态恢复(恢复层)、连接数与广播扇出(规模层)。延伸阅读可参考 https://plumephp.com/performance-load-testing/ 了解 k6/Locust 的负载编排与指标口径,https://plumephp.com/parallel-flaky-tests/ 了解如何治理实时测试在 CI 中的偶发失败,https://plumephp.com/e2e-testing-playwright/ 了解端到端场景里如何把实时交互纳入用户旅程验证。一句话收尾:实时系统的 bug 大多不在代码里,而在时间里——把时间、时序、规模这三样东西显式地测出来,实时通信才谈得上可靠。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。