面试知识库

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 全景流程图(文字版)#

1.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 0
python

为什么用 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()
plaintext

Step 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 消费

1.5 技术选型决策表#

选了什么备选方案为什么选它什么情况下换
Redis StreamsKafka / RabbitMQ / SQS项目已依赖 Redis(限流/执行槽/心跳),不引入新有状态服务;吞吐瓶颈在 LLM 不在队列告警量到每秒几千条或需跨数据中心 → Kafka
ZSET + Lua 做执行槽asyncio.Semaphore / RedlockSemaphore 进程内,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% 入队压测报告
Webhook500 请求 / 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 的优势是精细控制,劣势是需要主动调用(轮询)。


六、诚实边界#

  1. 3 Worker / 2 槽是本地开发配置,不代表生产能力。生产环境的参数需要根据 LLM 配额、基础设施能力和告警量实测后调整。
  2. 队列没有消息回溯/重放能力。Redis Streams 的消息被 ACK + XTRIM 后就没了。如果需要回放历史告警(评测场景),要从 Postgres 的 alerts 表重建——队列不是事实存储。
  3. DLQ 没有自动处理机制。DLQ 里的任务需要人工排查——当前没有工具、没有 dashboard、没有自动重试逻辑。DLQ 是”保证不丢”的兜底,但”丢进去之后怎么办”没做。
  4. 心跳只判存活不判健康。Worker 的心跳只说明进程活着,不说明它在正常工作——比如 Worker 进程活着但 event loop 死锁了,心跳还在刷但任务不推进。应该加”最近 N 秒有没有完成过任务”的活跃度判定。
  5. 优先级只有入队时静态分配,不支持运行时调整。如果一个 normal 任务排了很久变成了 critical(SLA 快超时了),不会自动提升优先级。