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 会被阻塞,导致:

三种处理:

策略 做法 适用
有界队列 + 丢弃 每个连接一个 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. 服务端下发 retry: 5000,并加随机抖动(客户端侧也要抖);
  2. 分批重启(每次 10% 实例);
  3. 重连请求也走限流(第 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 客户端重连次数、错误率

自测题

  1. 长连接系统的容量单位为什么从 QPS 变成并发连接数?怎么估算稳态连接数?
  2. 每条 SSE 连接占用哪些资源?其中最容易被忽略、最容易打爆下游的是哪一项?
  3. 慢客户端会造成什么问题?为什么"队列满了就断开"反而比"无限堆积"更安全?
  4. 多实例下为什么不需要会话亲和?什么条件成立这一点才成立?
  5. 重连风暴是怎么产生的?三种削峰手段是什么?

进入 keel 阅读