KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
06 · 连接治理与容量:长连接是一种新的容量单位 — keel 龙骨
这一章回答:1 万条 SSE 连接到底占多少资源,什么时候该拒绝新连接,以及多实例下怎么广播。
这一章回答:1 万条 SSE 连接到底占多少资源,什么时候该拒绝新连接,以及多实例下怎么广播。
传统接口用 QPS 衡量容量。长连接系统里,并发连接数成了新的容量单位——它不随请求结束而释放。这一章把它量化。
一、从 QPS 到并发连接数
稳态并发连接数 ≈ 同时在线的流数量
≈ 日活 × 同时使用比例 × 平均单流时长 / 平均使用间隔
举例:1 万日活、20% 同时在线、平均流时长 3 分钟 → 稳态约 2000 条连接。峰值要按 2~3 倍算。
每条连接的资源账单(异步框架下):
| 资源 | 每条占用 | 5000 条时 |
|---|---|---|
| 文件描述符 | 2(下游 + 上游) | 1 万 → 必须调 ulimit -n |
| 协程/任务 | 1~2 | 1 万协程(内存约几十~几百 MB) |
| nginx 连接 | 2 | 1 万 → 调 worker_connections |
| 应用内存 | 10~100KB(缓冲 + 上下文) | 50MB ~ 500MB |
| 事件源订阅 | 1(Redis 连接是否被独占是关键) | 见下节 |
⚠️ Redis 订阅是最容易被忽略的一项:如果每条 SSE 连接都独占一个 Redis 连接做 XREAD BLOCK,1 万条 SSE = 1 万个 Redis 连接,直接打爆 Redis。
正确做法:共享一个(或几个)轮询任务,把事件分发给本地的所有订阅者:
┌─ 本地订阅者 1 ─→ SSE 连接 1
一个 XREAD 轮询 ──┼─ 本地订阅者 2 ─→ SSE 连接 2
(每实例一个) └─ 本地订阅者 3 ─→ SSE 连接 3
即:Redis 侧的连接数与 SSE 连接数解耦,只与实例数相关。
二、背压:慢客户端会拖垮服务端
如果客户端读取很慢(弱网、移动端),服务端的 yield 会被阻塞,导致:
- 该协程卡住,占用内存(未发送的缓冲累积);
- 事件源的数据堆积;
- 极端情况下 OOM。
三种处理:
| 策略 | 做法 | 适用 |
|---|---|---|
| 有界队列 + 丢弃 | 每个连接一个 asyncio.Queue(maxsize=N),满了丢最旧的 |
允许丢中间态(token 流) |
| 有界队列 + 断连 | 满了就关闭连接(客户端会重连并从断点续) | 有事实日志时最稳 |
| 无界 + 监控告警 | 记录队列长度,超阈值告警 | 不推荐作为唯一手段 |
q: asyncio.Queue = asyncio.Queue(maxsize=200)
producer_task = asyncio.create_task(fill(q, run_id))
try:
while True:
evt = await q.get()
yield format(evt)
finally:
producer_task.cancel()
推荐"满了就断":因为有事实日志(第 04 章),断开后客户端会带着 Last-Event-ID 回来,丢的是连接不是数据。这比在服务端无限堆积内存安全得多。
三、准入:什么时候拒绝新连接
长连接系统必须有明确的准入,否则会被"连接数"这个新瓶颈打穿。
class StreamLimiter:
def __init__(self, max_global: int, max_per_user: int):
...
async def acquire(self, user_id: str) -> bool:
if self.global_count >= self.max_global: return False
if self.per_user.get(user_id, 0) >= self.max_per_user: return False
...
三层配额:
| 层级 | 限制 | 目的 |
|---|---|---|
| 全局 | 实例可承载的连接上限(如 5000) | 保护进程 |
| 单用户 | 同时最多 1~3 条 | 防误开/防刷 |
| 单资源 | 同一 run_id 的订阅数上限 | 防止一个热门资源被大量订阅 |
拒绝时返回明确的响应:
HTTP/1.1 429 Too Many Requests
Retry-After: 5
⚠️ 不要用 200 + 空流来"假装成功",客户端会以为连接正常而一直等。
四、多实例下的广播
事件源(Redis Stream / Kafka)本身就是共享的,所以"多实例都能读到同一份事实"是天然的。但要注意两种模式:
| 模式 | 实现 | 说明 |
|---|---|---|
| 按资源消费(常用) | 每个实例独立 XREAD 同一个 stream |
每个实例各自读全量,各自推给自己的订阅者(重复读但简单可靠) |
| Pub/Sub 扇出 | Redis SUBSCRIBE 频道,消息广播到所有实例 |
省去重读,但 Pub/Sub 不持久化,断线期间的消息会丢(需配合 Stream 做续播) |
推荐组合:Stream 做事实日志(负责续播),Pub/Sub 做实时通知(负责低延迟)。收到通知后去 Stream 里读新事件——这样既实时又能续。
# 伪代码:通知 + 续播
await pubsub.subscribe(f"run:{run_id}:notify")
async for _ in pubsub.listen():
entries = await redis.xread({key: last_id}, count=100, block=0)
for e in entries: yield format(e)
五、亲和性与负载均衡
SSE 不需要会话亲和(因为续播靠事实日志,任何实例都能接)。但发布时仍要注意:
- 优雅下线:新连接不再分发到旧实例,旧实例等现有流结束或超时;
- 滚动节奏:一次只下线一小部分实例,避免大量连接同时重连造成尖峰;
- 重连风暴:如果一次重启导致 1 万连接同时重连,会在几秒内产生 1 万个请求 + 1 万次回放读取。
削峰重连的三个手段:
- 服务端下发
retry: 5000,并加随机抖动(客户端侧也要抖); - 分批重启(每次 10% 实例);
- 重连请求也走限流(第 03 节的
StreamLimiter)。
// 客户端侧抖动(原生 EventSource 无法控制,需自己实现或用 fetch 方案)
await sleep(1000 + Math.random() * 3000);
六、监控指标(长连接系统的仪表盘)
| 指标 | 说明 | 告警方向 |
|---|---|---|
| 当前连接数 | 全局/每实例 | 接近上限 80% |
| 连接建立/断开速率 | 每秒 | 断开速率突增 = 网络或实例问题 |
| 重连次数 / 连接 | 平均每个流重连几次 | 突增说明不稳定 |
| 首字节时间(TTFB) | 连接建立到第一个事件 | 突增说明被缓冲了 |
| 事件发送速率 | 每秒事件数 | 突降为 0 而连接数不变 = 卡住 |
| 断开原因分类 | 客户端/超时/背压/错误 | 背压断开过多说明客户端太慢 |
| 队列堆积长度 | 每个连接的待发队列 | 持续增长 = 背压失效 |
最有价值的一条:事件发送速率 = 0 且 连接数 > 0。这是"连接还在但流已经死了"的唯一可靠信号,正是缓冲/生成器卡死的表现。
动手:可观察结果
| 产出 | 判断标准 |
|---|---|
| 一份容量测算 | 稳态连接数、单连接资源账单、峰值规划(含 fd / nginx / 内存) |
| 背压实现 | 有界队列 + 满则断;压测中人为制造慢客户端,服务端内存不失控 |
| 准入实现 | 全局 + 单用户 + 单资源三层配额;超限返回 429 而非空流 |
| 共享订阅改造 | 1 万 SSE 连接下 Redis 连接数不随之增长(验证解耦) |
| 重连风暴演练 | 一次重启 30% 实例,观察重连尖峰是否被抖动与限流削平 |
完成标志:给定"1 万同时在线"的目标,你能给出实例数、每实例连接上限、fd/内存配置、限流阈值四个数字,并说明超限时系统的表现。
故障注入
| 注入方式 | 观察 |
|---|---|
| 每条 SSE 独占一个 Redis 连接 | 5000 连接时 Redis 是否被打爆 |
| 客户端故意慢读(限速 1KB/s) | 服务端队列是否堆积、内存是否上涨 |
| 不设连接上限持续建连 | fd/内存是否耗尽,错误表现是什么 |
| 一次重启全部实例 | 是否出现重连风暴(Redis 读放大) |
下发固定 retry 无抖动 |
重连是否同步化形成尖峰 |
| 关闭优雅下线直接 kill | 客户端重连次数、错误率 |
自测题
- 长连接系统的容量单位为什么从 QPS 变成并发连接数?怎么估算稳态连接数?
- 每条 SSE 连接占用哪些资源?其中最容易被忽略、最容易打爆下游的是哪一项?
- 慢客户端会造成什么问题?为什么"队列满了就断开"反而比"无限堆积"更安全?
- 多实例下为什么不需要会话亲和?什么条件成立这一点才成立?
- 重连风暴是怎么产生的?三种削峰手段是什么?