M8 · 后台队列、Worker 与并发控制#
一句话定位: 把”接收请求”和”执行诊断”拆成两件事——API 只做校验入队立即返回,Worker 从 Redis Streams 消费任务、抢全局执行槽、跑诊断图,实现”接入量可以大、执行量严格可控”的削峰架构。
开场钩子#
场景#
V2 时代所有诊断都在 API 进程里同步跑——用户点”诊断”,FastAPI 开一个协程调 LLM,跑完再返回结果。三个人同时诊断就相安无事,10 个人同时诊断 API 就开始卡——LLM 的 RPM 限流打满,Postgres 连接池耗尽,一个慢诊断拖垮所有请求。更要命的是 Alertmanager 推告警——一场告警风暴 30 秒内推来 200 条 webhook,API 进程全被 LLM 调用占住,健康检查超时,Prometheus 误判服务挂了,推出新一轮告警——告警引发告警。
本质问题是请求速率和诊断执行速率耦合了。200 条告警进来不意味着需要同时跑 200 个诊断——经过 L2 降噪后可能只需要 3 个。但如果不拆开,告警入口的压力直接传导到 LLM 出站。
V3 的解法是经典的”队列 + Worker”架构:API 只做校验、归一化、入队,立即返回 task_id(webhook 202,< 50ms)。Worker 从 Redis Streams 消费任务,抢全局执行槽后才启动诊断。请求数量、队列深度、真实执行并发——三件不同的事,每一层各自管各自的上限。
面试官切入#
“你提到了全局执行槽——多 Worker 怎么共享并发上限?“
一、模块运作流程#
1.1 一句话定位#
本模块把 AIOps 系统的”告警入口”和”诊断执行”解耦——API 做接入削峰(入队即返),Worker 做执行限流(全局执行槽控总量),Redis Streams 做缓冲排序(四级优先级 + 死信兜底),三层协作保证”短时间收到的告警再多也不会打爆 LLM”。
1.2 全景流程图(文字版)#
并发请求 (可以很大)
│
▼
API 层 — 校验 + 归一 + 入队(Webhook 202 < 50ms)
│ ┌── critical
▼ ├── high
Redis Streams 四级优先级队列 ──────► ├── normal
│ └── low
│ 排队任务数 (按需积压)
▼
Worker 主循环:
① distributed_slot(wait=False) — 抢全局执行槽
│ ├── 抢到 → 继续
│ └── 没抢到 → 退避 0.5s,重试
② claim_stale_tasks — XAUTOCLAIM 回收 stale
③ read_tasks — Phase 1 非阻塞按优先级扫 → Phase 2 阻塞等待
④ run_diagnosis_graph — 执行诊断
⑤ xack — 消费完成
│
│ 真实诊断并发 (严格可控: 全局执行槽上限)
▼
结果落库 (Postgres)plaintext1.3 分步详解#
Step 1: 四级优先级队列
任务按 severity 分流到四条独立 Stream:{base}:critical / :high / :normal / :low。映射规则:critical/page/P0 → critical,high/P1 → high,info/low/P3 → low,其余 → normal。
两阶段消费:
- Phase 1 非阻塞扫:按 critical → high → normal → low → base 顺序逐条 stream 做
XREADGROUP(block=None),命中即返回——实现严格插队。 - Phase 2 阻塞等待:全空时对所有 stream 做一次
XREADGROUP(block=block_ms),有新消息即醒来,下一轮重新按优先级取。
为什么 Phase 1 不一次性读所有 stream?因为 Redis Streams 的 XREADGROUP 对多 stream 返回的是”每条 stream 最多 count 条”——如果 normal 里积了 100 条而 critical 来了 1 条,一次读不保证 critical 先被消费。逐条 stream 扫保证严格优先级。
Step 2: 全局执行槽(Redis ZSET + Lua)
# 数据结构: ZSET aiops:limiter:worker_diagnosis
# member = "host:pid:uuid" (唯一 token)
# score = 过期时间 (ms, 服务端时间)
# Lua 原子抢槽:
ZREMRANGEBYSCORE(0, now) # 清过期
count = ZCARD # 数当前
if count < max_slots:
ZADD(token, now + ttl) # 占位
return 1
return 0python为什么用 Lua:check-then-act 必须原子——先 ZCARD 判断没满再 ZADD 占位,不用 Lua 会有竞态。
心跳续期:长诊断(30-60s)超过默认 TTL 会被误清。_REFRESH_LUA 只在 token 还在时更新 score。
SlotHandle.pause/resume:等人工审批时先释放槽(ZREM),审批完再重抢——审批期间不占诊断资源。通过 ContextVar[SlotHandle] 暴露给 tool_runner。
Fail-Open:Redis 不可用时 acquire 返回 __fail_open__ token,直接放行。理由:队列本就强依赖 Redis,真挂了整条链路都不工作——限流层不叠硬失败。
Step 3: Worker 主循环
关键设计:先抢槽再读 stream。抢不到槽就不碰 stream——不会把消息从 stream 里取出来堆在 PEL(Pending Entry List)。
while running:
slot = distributed_slot(wait=False)
if not slot:
await sleep(0.5) # 退避
continue
task = claim_stale_tasks() or read_tasks()
if task:
run_diagnosis_graph(task)
xack(task)
slot.release()plaintextStep 4: XAUTOCLAIM 崩溃恢复
Worker 崩溃后,它领走但未 ACK 的消息会留在 PEL。claim_stale_tasks 按优先级顺序在每条 stream 上调 XAUTOCLAIM,把空闲超过 reclaim_idle_ms(默认 15 分钟,= 诊断超时 10 分钟 + 5 分钟安全边际)的 pending 消息转给当前 Worker。
关键区分:刚被领走正在执行的消息 idle 很短,不会被误抢。只有真正 stale 的(Worker 崩溃了 15 分钟没心跳)才会被回收。
Step 5: 死信队列(DLQ)
入 DLQ 条件:(a) 消息缺 task_id;(b) task 在 Postgres 找不到(脏数据);(c) attempts ≥ max_attempts(默认 3)。
DLQ 操作:原始消息 + 失败原因写入 DLQ stream(xadd),然后在原消息所在的优先级 stream 上 xack——确保原 stream 不积压已知无法处理的消息。
Step 6: 心跳
Worker 写一个带 TTL(默认 30s)的 Redis key aiops:worker:{group}:{name}:heartbeat,每 10s 刷新。status() 方法通过 key 是否存在判断 Worker alive。30s 没刷新 = Worker 挂了。
Step 7: 可观测性指标
status() 聚合返回:每条 stream 的 XLEN / lag / pending,DLQ 深度,Worker 存活列表,执行槽占用/总量。关键定义:backlog = lag(未投递) + pending(已投递未 ACK),不是 XLEN(ACK 后不立即下降)。
1.4 数据流 trace#
场景:一条 critical 告警入队到被 Worker 消费
输入: Alertmanager webhook POST {
alerts: [{
status: "firing",
labels: {alertname: "MySQLDown", service: "prod-mysql", severity: "critical"},
...
}]
}
Step 1 — API 入队:
归一化 → NormalizedAlert
level_for_severity("critical") → "critical"
xadd aiops:diagnosis:critical {task_id: "task-123", severity: "critical", ...}
HTTP 202 + {task_id: "task-123"} ← < 50ms
Step 2 — Worker 抢槽:
EVAL _ACQUIRE_LUA aiops:limiter:worker_diagnosis
→ ZCARD=1, max=32 → 占位成功
Step 3 — Worker 读任务:
Phase 1: XREADGROUP aiops:diagnosis:critical → 命中 task-123
(不需要看 high/normal/low——critical 已有消息)
Step 4 — 执行诊断:
run_diagnosis_graph(task-123)
→ L2 预处理 → Tier 1 miss → Tier 2 miss → Tier 3 hypothesis-driven → RCA confirmed
Step 5 — 落库 + ACK:
mark_task_succeeded(task-123)
XACK aiops:diagnosis:critical task-123
Step 6 — 释放槽:
ZREM aiops:limiter:worker_diagnosis token-xxxplaintext1.5 技术选型决策表#
| 选了什么 | 备选方案 | 为什么选它 | 什么情况下换 |
|---|---|---|---|
| Redis Streams | Kafka / RabbitMQ / SQS | 项目已依赖 Redis(限流/执行槽/心跳),不引入新有状态服务;吞吐瓶颈在 LLM 不在队列 | 告警量到每秒几千条或需跨数据中心 → Kafka |
| ZSET + Lua 做执行槽 | asyncio.Semaphore / Redlock | Semaphore 进程内,Redlock 是互斥锁不是计数器 | N/A(ZSET 是正确的数据结构选择) |
| 四级优先级 | 单队列 FIFO / 更多级 | 对标 PagerDuty 的 critical/high/low 分级,4 级覆盖 Alertmanager 常见 severity | 实测发现大部分告警集中在 normal → 简化为 2 级 |
1.6 口述脚本#
开场(20s):“AIOps 系统面临的并发问题不是常规 web 应用的高 QPS——而是告警风暴。30 秒内 200 条 webhook,如果全部同步跑诊断,LLM 的 RPM 先打满,Postgres 连接池再耗尽。”
架构(30s):“解法是把请求数、排队任务数、真实执行数拆成三件事——API 只做入队(< 50ms 返回),Worker 从 Redis Streams 消费,全局执行槽控制总并发。200 条告警进来,经 L2 降噪后可能只需要 3 个诊断——队列吸收了差值。”
执行槽重点(20s):“执行槽用 Redis ZSET + Lua 原子抢占,多 Worker 共享同一个上限。长诊断有心跳续期,审批等待时释放槽让出资源。”
留钩子(10s):“最有意思的是 Fail-Open 设计——Redis 挂了限流器怎么办。“
二、踩坑实录#
坑 1:XLEN 不是真实 backlog#
- 现象:前端展示队列深度用
XLEN,但任务被 Worker 消费并 ACK 后 XLEN 没有立即下降——因为 Redis Streams 的XLEN返回的是 stream 里的总消息数,ACK 不删消息(需要XTRIM)。运维看到”队列深度 200”但实际早处理完了。 - 修法:backlog 的正确定义是
lag + pending。lag = 未投递(stream 里还没被任何 consumer 读到的消息数,XINFO GROUPS的 lag 字段)。pending = 已投递未 ACK(在 PEL 里的消息数)。两者之和才是”还没处理完的任务数”。 - 教训:Redis Streams 的 XLEN 语义和你以为的不一样。它是 stream 的物理长度,不是逻辑队列深度。这是 Redis Streams 与传统消息队列的显著差异。
坑 2:先读 stream 再抢槽导致 PEL 堆积#
- 现象:早期 Worker 是”先读消息再抢槽”——消息取出来了但没槽跑,消息留在 PEL 里。另一个 Worker 的 XAUTOCLAIM 把它回收走,第一个 Worker 拿到槽后发现消息没了,空转一次。
- 修法:改为”先抢槽再读 stream”——抢不到就不碰 stream,消息安静地留在 stream 里等有容量的 Worker 来取。
- 教训:消费和处理的顺序很重要。在有限资源场景下,应该先确认有资源再拿任务,而不是先拿任务再等资源。
坑 3:审批等待占着执行槽#
- 现象:处置走飞书审批时,Worker 在
interrupt()处等待。等待期间执行槽被占着——一个审批等 5 分钟,5 分钟里少了一个诊断资源。3 个 Worker 2 个槽,一个在等审批,实际可用只剩 1 个。 - 修法:
SlotHandle.pause()——等审批时释放槽(ZREM),审批完resume()再重抢。通过ContextVar[SlotHandle]暴露给 tool_runner,不需要修改图的逻辑。 - 教训:长时等待 ≠ 长时执行。等待外部输入(审批/验证观察期)时应该释放计算资源,而不是一直占着。
坑 4:重试消息的状态和队列不同步#
- 现象:任务失败后先
mark_task_retry_pending(DB 改状态为 pending),再xadd(重新入队),再xack(ACK 原消息)。如果 xadd 后 Worker 崩了——DB 状态是 pending,新消息在 stream 里,但旧消息没 ACK(还在 PEL)。下次回收会拿到旧消息,而新消息也会被正常消费——同一个任务被执行两次。 - 修法:顺序改为”先 xadd 再 xack”——即使 xack 前崩了,旧消息被回收后通过幂等检查(task 已 succeeded)直接跳过。幂等比顺序更重要。
- 教训:分布式操作的顺序和幂等都很重要。如果无法保证原子性,至少保证幂等——让重复执行不产生副作用。
三、量化评估#
3.1 评估数据集构建#
- 压测数据:
scripts/loadtest.py,200 请求 / 100 并发提交后台诊断,500 请求 / 100 并发 webhook。数据来源docs/PRESSURE_TEST_REPORT.md。 - 队列测试:
FakeRedisStreams复刻 Redis Streams 关键语义(xadd/xreadgroup/xack/xautoclaim + PEL + consumer group),10+ 个测试用例。
3.2 评估方法#
- 压测脚本:
python scripts/loadtest.py --scenario all - 对照组:无队列(API 直接跑诊断)vs 有队列(异步入队 + Worker 消费)
- 可复现:本地 3 Worker / 2 执行槽 / Redis Streams / Postgres
3.3 结果与解读#
| 指标 | 值 | 来源 |
|---|---|---|
| 读接口 | 3000 请求 / 200 并发,100% 成功 | 压测报告 |
| 后台诊断提交 | 200 请求 / 100 并发,100% 入队 | 压测报告 |
| Webhook | 500 请求 / 100 并发,告警全部接收归一落库 | 压测报告 |
| 全局执行槽 | 始终不超过 2/2(3 Worker 共享) | 压测报告 |
| 优先级插队 | critical 在 low 之前被消费 ✅ | test_incident_queue.py |
| 崩溃恢复 | XAUTOCLAIM 回收 stale 消息 ✅ | test_incident_queue.py |
局限性:3 Worker / 2 槽是本地开发配置,与生产环境差距大。真正的压力在 LLM RPM/TPM,不在队列吞吐——但这正是分层治理的意义:队列不需要快,它只需要把压力从 API 层卸下来。
四、面试问答#
基础题(必问级)#
Q1: 为什么不用 Kafka?Redis Streams 不是更简单的消息队列吗?#
答: 两个原因。一是项目已经依赖 Redis(限流、执行槽、心跳、冷却期),引入 Kafka 增加一个独立的有状态服务,运维成本上升但收益不大。二是瓶颈不在队列吞吐——一个诊断任务跑 30-60 秒 LLM 调用,队列每分钟处理几十个任务就够了,Redis Streams 绰绰有余。Kafka 的优势在海量吞吐(百万级 TPS)、跨数据中心复制和流式计算集成,这些当前用不上。如果告警量到了每秒几千条,或者需要多机房部署,那时候换 Kafka 才有实际收益。
追问: Redis Streams 的持久化不如 Kafka 吧?消息丢了怎么办? 答: Redis 的持久化靠 RDB + AOF,确实不如 Kafka 的分区副本。但消息丢失的影响有限——告警源(Alertmanager)有重推机制,丢了的消息会被重推。更重要的是,任务的状态不只在 Redis 里——入队时同时落 Postgres(
diagnosis_tasks表),即使 Redis 数据全丢,可以从 Postgres 重建 pending 任务。Redis 是”加速通道”不是”唯一真相”。
Q2: 全局执行槽为什么不用 asyncio.Semaphore?#
答: asyncio.Semaphore 是进程内的——3 个 Worker 进程各有各的 Semaphore,无法共享上限。全局执行槽的需求是”3 个 Worker 共享 2 个真实执行名额”,必须用进程间共享的数据结构。Redis ZSET + Lua 是正确的选择:ZSET 的 ZCARD 就是计数器,Lua 保证原子性,score 存过期时间实现自动清理。
追问: 为什么不用 Redlock? 答: Redlock 是分布式互斥锁(最多 1 个持有者),不是计数信号量(最多 N 个持有者)。全局执行槽需要”最多 2 个并发”,不是”最多 1 个并发”。ZSET 天然是计数器——ZCARD 返回当前持有者数量,ZADD 新增持有者,ZREM 释放。
进阶题(区分度)#
Q3: Fail-Open 设计——Redis 挂了所有请求都不限流,LLM 不会被打爆吗?#
答: 会有风险,但 Fail-Open 是比 Fail-Closed 更好的选择。理由:队列本身强依赖 Redis——Redis 挂了消息都读不出来,Worker 的整条链路都不工作。限流层在这时候硬失败(Fail-Closed)等于”告诉还没进入队列的请求也别进来了”——但这些请求是 API 层的,本来就不直接跑 LLM。Fail-Open 影响的是”此刻有多少个已经在跑的诊断”——如果 Redis 只是闪断几秒,可能会有短暂超限,但 LLM 自身的 RPM 限流会兜底。如果 Redis 持续挂了,整条链路都挂了——限流不限流没区别。
Q4: 15 分钟回收阈值——如果诊断真的跑了 14 分钟呢?#
答: 14 分钟的诊断不会被误回收——因为 Worker 有心跳续期。正常运行的 Worker 每 10s 刷新消息的 idle 计时(通过读取消息后的 ACK 时间窗口)。只有 Worker 真的崩了(15 分钟没心跳)才会 idle 超过阈值。15 分钟 = 诊断超时 10 分钟 + 5 分钟安全边际。如果诊断超过 10 分钟还没完,AgentHarness 会强制终止(步骤上限 / token 上限),不存在”跑 14 分钟”的合法场景。
压力题(面试官挑战设计决策)#
Q5: 本地 3 Worker / 2 槽——跟生产差距大不大?你怎么证明这套东西上了生产能用?#
答: 差距肯定有。本地压测验证的是机制正确性——队列削峰、执行槽限制、优先级抢占、DLQ 这些逻辑在并发下是否按预期工作。生产环境的变量更多:网络延迟、Redis/Postgres 集群化、LLM API 的 RPM 波动。当前配置(3 Worker / 2 槽)是本地开发环境的经验值,上生产需要根据 LLM 配额和基础设施能力重新调参。但核心架构——请求接入和诊断执行分离、全局执行槽跨进程共享——这个拓扑不需要改,改的是参数。
追问: 那执行槽用 Redis ZSET + Lua,Redis 本身是单线程的,争抢很频繁时会不会成为瓶颈? 答: 当前场景不会。Lua 脚本只做 ZREMRANGEBYSCORE + ZCARD + ZADD 三个操作,执行时间微秒级。执行槽的争抢频率取决于诊断任务的提交速率,不是请求速率——请求到执行之间有队列缓冲。即使 100 个 Worker 争抢 10 个槽位,每次 Lua 调用也就几微秒。真正的 Redis 瓶颈更可能出现在限流器的 INCR 高频调用上,但那也是简单的单键操作。
Q6: 这套并发设计的投入产出比怎么样?值得吗?#
答: 队列 + Worker + 执行槽 + 心跳 + DLQ 大概花了两周。如果只做 demo,直接在 API 进程里跑 Agent 就够了。但如果要回答”高并发怎么办、任务丢了怎么办、Worker 崩了怎么办”这些面试问题,就必须有真实的队列和 Worker。从面试角度,这是项目里最能体现系统设计能力的部分——面试官问的不是”你会不会用 Redis”,而是”你能不能设计一个从接入到执行到失败恢复的完整链路”。投入产出比很高。
五、前沿概念与延伸#
5.1 Backpressure(背压)#
是什么:当下游处理速度跟不上上游产生速度时,系统主动向上游施加”压力”(减缓发送速率、拒绝新请求、丢弃低优先级消息),而不是无限缓冲直到内存耗尽。
为什么要这么做:无限缓冲的系统在面对突发流量时会积压越来越多的未处理消息,最终内存溢出或延迟无限增大。背压让系统在压力下优雅降级而不是崩溃。
这么做的理由:相比无限制入队(内存终将耗尽)或硬拒绝(丢失告警),背压提供了中间路线——队列有上限,达到上限后新任务降级处理(低优先级丢弃、返回 429)。
举例说明:本项目的三层背压——(1) API 层限流(固定窗口 + Redis INCR,超限返回 429);(2) 队列容量(Redis Streams 可配 MAXLEN);(3) 全局执行槽(Worker 消费上限)。三层各自独立——API 限流保护队列不爆,执行槽保护 LLM 不被打满。
延伸问答:
面试官:“Reactive Streams 的背压和你的有什么不同?” 答:Reactive Streams(如 Project Reactor)的背压是端到端的——消费者通过
request(n)告诉生产者”我只能处理 n 条”。我们的背压是分层的——每层各自控制自己的上限,层间通过队列解耦。分层背压更适合异构系统(API、队列、Worker 是独立进程),端到端背压更适合进程内的流式处理。
5.2 Consumer Group 与 Exactly-Once 语义#
是什么:Consumer Group 是消息队列的消费者分组机制——同一条消息只会被组内一个消费者处理,实现负载均衡。Exactly-Once 是”每条消息恰好被处理一次”的语义保证。
为什么要这么做:多 Worker 消费同一个队列,如果没有 Consumer Group 就会重复消费。Exactly-Once 在分布式环境下很难实现——通常退化为 At-Least-Once + 幂等。
这么做的理由:Redis Streams 原生支持 Consumer Group(XREADGROUP),At-Least-Once 由 PEL + XACK 保证——消息被读出后如果没 ACK 就留在 PEL,崩溃恢复时 XAUTOCLAIM 回收。Exactly-Once 通过幂等键实现——相同 task_id 不会被二次执行。
举例说明:本项目的幂等保护:task.status == succeeded 时收到重复消息直接 ACK 跳过;IdempotencyGuard 用 Redis SET NX 首见放行复见拒绝。
延伸问答:
面试官:“Redis Streams 的 XAUTOCLAIM 和 Kafka 的 Consumer Group rebalance 有什么区别?” 答:Kafka rebalance 是全组级别的——一个消费者挂了,Coordinator 触发 rebalance 把它的分区分配给其他消费者。Redis Streams 的 XAUTOCLAIM 是消息级别的——只回收特定 idle 超过阈值的 pending 消息,不影响其他消息的消费。Kafka rebalance 的优势是完全自动,劣势是 rebalance 期间全组停止消费(stop-the-world)。XAUTOCLAIM 的优势是精细控制,劣势是需要主动调用(轮询)。
六、诚实边界#
- 3 Worker / 2 槽是本地开发配置,不代表生产能力。生产环境的参数需要根据 LLM 配额、基础设施能力和告警量实测后调整。
- 队列没有消息回溯/重放能力。Redis Streams 的消息被 ACK + XTRIM 后就没了。如果需要回放历史告警(评测场景),要从 Postgres 的
alerts表重建——队列不是事实存储。 - DLQ 没有自动处理机制。DLQ 里的任务需要人工排查——当前没有工具、没有 dashboard、没有自动重试逻辑。DLQ 是”保证不丢”的兜底,但”丢进去之后怎么办”没做。
- 心跳只判存活不判健康。Worker 的心跳只说明进程活着,不说明它在正常工作——比如 Worker 进程活着但 event loop 死锁了,心跳还在刷但任务不推进。应该加”最近 N 秒有没有完成过任务”的活跃度判定。
- 优先级只有入队时静态分配,不支持运行时调整。如果一个 normal 任务排了很久变成了 critical(SLA 快超时了),不会自动提升优先级。