面试知识库
高 困难

Agent长任务队列与服务化#

一句话答案#

Agent 一轮要跑几十秒到几分钟、按 token 花钱,不能挂在请求线程里同步跑:API 只做鉴权、配额预扣、幂等判定和入队,立刻返回 task_id;worker 从 Redis Stream 消费者组领任务执行,写完终态才 XACK;崩溃的留在 PEL,由 XAUTOCLAIM 接管,接管后是重跑还是直接给「中断」终态,要按成本和用户是否还在场来取舍。队列只提供 at-least-once,真正防重复的是业务层幂等(条件更新的状态机 + 按 run_id 幂等的结算)。

核心要点

1. 为什么不能在请求里同步跑#

问题同步跑的后果
超时链路Nginx proxy_read_timeout 默认 60 秒、云负载均衡也有空闲超时,Agent 还没跑完连接就被网关切断,模型调用却照样在花钱
发布重启滚动发布杀掉 API 进程,正在跑的任务全部丢失,没有任何记录
没有背压突发流量直接变成几百个并发 LLM 调用,打满下游限流后全部失败
无法分别扩容API 是轻量 IO,Agent 执行吃内存和下游配额,绑在一个进程里只能一起扩
重试放大客户端超时重试 = 同一个问题真跑两遍、扣两份钱

拆开后的形态:

flowchart LR
    C[客户端] -->|POST /task| A[API 进程]
    A -->|1 鉴权 2 预扣配额 3 幂等判定| DB[(MySQL)]
    A -->|4 XADD| S["Redis Stream"]
    A -->|返回 task_id| C
    S -->|XREADGROUP| W1[worker 1]
    S -->|XREADGROUP| W2[worker 2]
    W1 -->|事件 Pub/Sub| A
    A -->|WS / SSE 推送| C
    W1 -->|终态 + 结算| DB
    W1 -->|XACK| S

客户端靠 WebSocket / SSE 收进度,断线后按事件 ID 补发,最终结果可以查库拿到。API 进程里一个 Agent 都不跑。

2. Redis Stream 消费者组:at-least-once 的实现#

Stream 与 List / Pub/Sub 的对比见 Redis实现消息队列(List-PubSub-Stream),这里讲任务队列用到的完整链路:

命令作用关键点
XGROUP CREATE key g $ MKSTREAM建消费者组组已存在会报 BUSYGROUP,启动时捕获忽略
XADD key MAXLEN ~ N * f v入队~ 近似裁剪,防止 Stream 变大 Key
XREADGROUP GROUP g c COUNT n BLOCK ms STREAMS key >领新消息> 表示「从没投递给本组任何消费者的」;领走的消息进 PEL
XREADGROUP ... STREAMS key 0读自己 PEL 里的历史worker 重启后先处理自己名下没 ack 的
XACK key g id确认从 PEL 删除,写完终态后才调
XPENDING key g看 PEL 摘要只有总数、最小/最大 ID、每个消费者的条数
XPENDING key g [IDLE ms] - + count [consumer]看 PEL 明细扩展形式才能看到每条的空闲时间和投递次数
XAUTOCLAIM key g c min-idle start COUNT n认领别人名下超时未 ack 的消息Redis 6.2+;投递次数 +1(JUSTID 不加),超过阈值转死信;7.0 起返回三元素:下次扫描游标、认领到的消息、已从 Stream 删除而被清出 PEL 的 id
XREADGROUP ... CLAIM min-idle STREAMS key >一条命令里先认领空闲超时的 pending,再读新消息Redis 8.4+;认领到的条目额外带空闲时长和投递次数,可替代单独的 XAUTOCLAIM 循环
XNACK key g FAIL|SILENT|FATAL IDS n id...主动把消息放回 PEL,不等空闲超时即可被别人认领Redis 8.8+;SILENT 投递次数减 1(适合优雅退出)、FATAL 标记为永久失败(毒消息)
XACKDEL key g [KEEPREF|DELREF|ACKED] IDS n id...ack 和删除 Stream 条目一步完成Redis 8.2+;多个消费者组时用 ACKED 只删所有组都 ack 过的条目

长任务专有的坑:PEL 里的「空闲时间」从投递那一刻算起,任务跑着但没 ack 也在累加。min-idle 如果小于任务的最长执行时间,另一个 worker 会把正在跑的任务当成崩溃接管:按上面的写法会被误判成中断,如果接管后重跑,就是同一个任务跑两遍。两种解法:min-idle 设成大于单轮超时;或者执行中的 worker 定期对自己手里的消息 XCLAIM ... JUSTID 给自己续约(认领会重置空闲时间)。

3. 幂等:队列保证不了 exactly-once#

重复投递来源:worker 执行完、ack 前崩溃;XAUTOCLAIM 误领;客户端超时重发。通用手段见 消息重复与幂等方案,Agent 场景的三层:

  1. 入口去重:请求指纹 SET dedup:{hash} 1 NX EX <窗口秒数>,查和登记是一步原子操作。Redis 不可达时有两种取舍:直接报错(503),代价是可用性;或者退回进程内字典,代价是多副本下重复请求会漏过去。涉及扣费和写操作时一般选前者。
  2. 任务状态机条件更新:领任务时 UPDATE runs SET state='running', worker=? WHERE id=? AND state='queued',影响行数为 1 才执行;判定和占位在同一条语句,数据库保证并发两条只有一条成功。已经是 running 的任务被重投时不会再跑一遍,交给崩溃接管的策略处理(见第 5 节)。
  3. 副作用幂等:下单、发邮件这类工具调用带 operation_id,下游按它去重;结算按 run_id 幂等(见第 6 节)。

4. 和 Celery / Kafka / RabbitMQ 的取舍#

方案适合长任务要注意的点
Redis Stream已有 Redis、量不大、想少一个中间件数据在内存要 MAXLEN;min-idle 与任务时长的关系要自己管;没有内置延迟队列和死信
CeleryPython 生态、要现成的重试/定时/结果后端默认执行前 ack,长任务要 acks_late=True;Redis 做 broker 时有 visibility_timeout(官方文档写默认 1 小时),任务超过它没 ack 会被重投给别的 worker;原生 async 支持弱,跑 asyncio 代码要自己在任务里起事件循环
Kafka高吞吐事件流、要回放传统消费者组按分区提交 offset,一条慢任务会堵住整个分区;处理时间超过 max.poll.interval.ms(默认 5 分钟)触发 Rebalance。Kafka 4.2 起 share groups(KIP-932,Queues for Kafka)进入生产可用,支持逐条 ack、投递次数和续期,但对「一条消息跑几分钟」的场景仍要单独评估
RabbitMQ要逐条 ack、优先级、死信、延迟插件prefetch=1 防止一个 worker 囤任务;有消费确认超时(官方文档写默认 30 分钟,4.3 起只有 quorum queue 支持该超时),超长任务要调大

选型原则:Agent 任务是「低 QPS、单条很贵、执行很长」,要的是逐条 ack 和崩溃接管,而不是吞吐——Redis Stream 或 RabbitMQ 更贴合,Kafka 传统消费者组更适合事件日志和埋点。通用 MQ 对比见 MQ选型对比。

5. 取消、优雅退出、崩溃接管#

取消:还在队列里的任务靠按 task_id 打的取消标记拦截,已经在跑的靠 Pub/Sub 广播取消指令、持有者调 task.cancel()。跨实例怎么送达见 实时推送断线续传与多实例事件转发,CancelledError 的传播见 Python异步编程与asyncio事件循环。

SIGTERM 优雅退出:

async def run_worker():
    stop = asyncio.Event()
    loop = asyncio.get_running_loop()
    loop.add_signal_handler(signal.SIGTERM, stop.set)

    consumer = asyncio.create_task(worker_loop(r, NAME, stop))  # 领到的任务放进模块级 running
    await stop.wait()                     # 1. 收到 SIGTERM:不再领新任务
    consumer.cancel()                     # 只取消领取循环;running 里是独立 Task,不受影响
    if running:                           # asyncio.wait 传空集合会抛 ValueError
        done, pending = await asyncio.wait(set(running), timeout=GRACE_SECONDS)  # 2. 等在跑的跑完
        for t in pending:                 # 3. 排空超时:取消,handle 里写 interrupted 并 ack
            t.cancel()
        await asyncio.gather(*pending, return_exceptions=True)
python
  • GRACE_SECONDS 必须 ≥ 单轮 Agent 超时,否则每次发布都会主动掐掉一批正常任务;K8s 的 terminationGracePeriodSeconds 要再大于它(加上 preStop 和收尾时间),否则 SIGKILL 先到。Pod 终止流程见 Pod生命周期。
  • 排空超时的任务有两种处理:不 ack 留在 PEL 等别的 worker 重跑(Redis 8.8+ 可以用 XNACK ... SILENT 主动放回,不必等 min-idle),或者当场写 interrupted 终态并 ack、让用户手动重发。前者适合用户还在等、结果有价值的任务,代价是用户可能早已离开,模型调用却真的又花一次;后者多一次手动操作但成本可控。
  • 收尾顺序有讲究:先发中断事件 → 释放会话占位 → 写终态 → ack。如果先写终态、后释放占位,客户端看到终态马上重发,会被还没释放的占位挡成「已在运行」。

崩溃接管(kill -9 / OOM): 来不及收尾,消息留在 PEL。空闲时间超过 min-idle(大于单轮超时,相当于租约)后,reclaim_loop 用 XAUTOCLAIM 领回。领回后的处理同样是取舍:

  • 重跑:把 running 重置回 queued 再交给 handle,适合短任务、结果仍有人等的场景;要配合投递次数上限,并保证工具副作用幂等。
  • 给中断定论:按上面的收尾顺序把 queued / running 条件更新成 interrupted 再 ack,由用户决定是否重发,适合单轮很贵、跑得很久的 Agent 任务。

数据库里的会话占位和预扣行另带 expires_at,接管方还没处理到时也会过期失效,不会一直占着额度。

6. 任务状态机、结果持久化、预扣-结算#

queued ──领取成功──▶ running ──▶ succeeded
   │                  ├────▶ failed       (可重发)
   │                  ├────▶ cancelled    (用户取消)
   │                  └────▶ interrupted  (发布/排空超时被掐)
   └──超时无人领取──▶ failed
plaintext
  • 每次状态迁移都用条件更新 WHERE state IN (...),非法迁移影响行数为 0,天然幂等。终态集合在代码里单点定义,前后端共用。
  • 结果持久化:最终答案和产物写 DB / 对象存储;执行中的事件写一条按会话分的 Stream(带 MAXLEN 和 TTL),客户端断线重连带 last_event_id 补发。worker 和 API 读写的产物路径要从同一个配置派生,否则 worker 写了 API 读不到。

预扣-结算式计费(配额体系见 多租户隔离与配额治理):事后记账挡不住并发透支——N 个请求同时进门都读到「额度够」,跑完一起扣就超了。改成进门预扣:

-- 进门(同一事务):锁用户行 → 数在飞的 run → 扣预估额度 → 插一行 hold
SELECT ... FROM users WHERE id = ? FOR UPDATE;
INSERT INTO quota_holds(run_id, user_id, amount, state, expires_at)
VALUES (?, ?, ?, 'active', NOW() + INTERVAL ? SECOND);   -- 过期时间略大于单轮超时

-- 结束时结算:条件更新就是幂等判据,重复投递第二次影响 0 行
UPDATE quota_holds SET state = 'settled', actual = ? WHERE run_id = ? AND state = 'active';
-- 仅当上一句影响 1 行,才在同一事务里记账
sql
  • 结算和账本累加在同一个事务里,重投不会双记。
  • 被幂等拒绝的请求(重复、已在运行)要把预扣还回去,每条拒绝路径都要有对应的释放。
  • 取消、中断、超时这些终态也要走结算(按实际花费),不能只在成功路径上结算。

面试回答(2分钟版)

Agent 一轮要几十秒到几分钟,而且每步都在花钱,放在请求线程里同步跑会被网关超时切断,发布一重启任务就丢,突发流量也没有背压,客户端重试还会让同一个问题真跑两遍。所以常见做法是拆成 API 和 worker 两个进程:API 只做鉴权、配额预扣、幂等判定,然后 XADD 进 Redis Stream,立刻返回 task_id,进度通过 WebSocket 或 SSE 推。worker 用 XREADGROUP 在消费者组里领任务,领走的消息进 PEL,写完终态才 XACK;worker 崩溃后,消息留在 PEL,别的 worker 用 XAUTOCLAIM 接管,接管后可以重跑,也可以写中断终态再 ack 让用户重发,Agent 这种单轮很贵的任务通常选后者。这里有个长任务特有的坑:min-idle 要大于单轮超时,否则正在跑的任务会被当成崩溃接管。队列只保证 at-least-once,防重复靠业务层:入口用 Redis SET NX EX 做请求指纹,领任务用状态机条件更新,只有影响行数为 1 才执行,结算按 run_id 条件更新保证不双记。优雅退出是 SIGTERM 后停止领新任务,等待时间要不小于单轮超时,排空超时的任务要么放回队列重跑、要么写中断终态再 ack,收尾顺序是先释放占位再写终态。选型上 Kafka 传统消费者组按分区提交 offset,一条慢任务会堵分区,不适合;Celery 要开 acks_late 并注意 visibility_timeout。结合项目时可以讲:额度在哪一层预扣、同会话唯一运行靠什么原子操作保证、被掐断的任务选择重跑还是给定论,以及用什么数据验证这个取舍。

追问与易错

追问方向:

  • “XREADGROUP 里 > 和 0 有什么区别?” → > 只领本组从没投递过的新消息;0 返回的是当前消费者自己 PEL 里已投递未 ack 的历史消息,worker 重启后先用 0 处理自己名下的遗留,再切到 >。
  • “worker 执行完、XACK 之前崩溃了会怎样?” → 消息留在 PEL,空闲时间超过 min-idle 后被别的 worker 用 XAUTOCLAIM 领回。任务如果已经写了终态,条件更新影响 0 行,直接 ack;如果还停在 queued / running,按接管策略处理:选择「给定论」就条件更新成 interrupted 再 ack、用户手动重发;选择「重跑」就把状态重置回 queued 再执行,并检查投递次数。claim_run 只认 queued,所以没经过接管策略的重投消息不会被当成新任务再跑一遍。
  • “XAUTOCLAIM 的 min-idle 怎么定?” → 必须大于单任务最长执行时间(通常取单轮超时加一段余量),否则会抢走正在跑的任务;更精细的做法是执行中的 worker 定期对自己的消息 XCLAIM JUSTID 续约。
  • “一条任务反复失败怎么办?” → 如果失败就写 failed 并 ack、接管也不重跑,就不会反复投递。如果选择接管后重跑:XPENDING 扩展形式能看到投递次数(Redis 8.4+ 的 XREADGROUP ... CLAIM 会直接在回复里带上),领回时检查,超过阈值(如 3 次)就写入死信 Stream 并 XACK 原消息,同时把任务标 failed 通知用户,别让毒消息一直循环;Redis 8.8+ 也可以用 XNACK ... FATAL 标记永久失败。
  • “为什么 Agent 任务不用 Kafka?” → Kafka 按分区顺序提交 offset,一条跑几分钟的任务会堵住整个分区;处理时间超过 max.poll.interval.ms 还会触发 Rebalance 把分区转走,导致重复消费。
  • “SIGTERM 后等多久?K8s 怎么配?” → worker 排空等待时间 ≥ 单轮 Agent 超时;K8s 的 terminationGracePeriodSeconds ≥ preStop + 排空等待 + 收尾余量,否则 SIGKILL 先到,收尾逻辑跑不完。
  • “正在跑的任务怎么取消?” → API 不知道任务在哪个 worker,就通过 Pub/Sub 广播取消指令,持有者调 task.cancel(),在捕获 CancelledError 的分支里写 cancelled 终态并结算;还没被领走的任务靠按 task_id 的取消标记拦截。
  • “预扣和事后记账比好在哪?预扣的钱没还怎么办?” → 预扣在进门时锁用户行扣额度,并发请求串行判断,不会透支;worker 被 kill -9 时来不及结算,hold 行带 expires_at,过期后不再计入占用。
  • “配额预扣和幂等判定谁先谁后?” → 两种排法各有代价。一种取舍是先判幂等再预扣:重复请求不占额度,但判定和入队之间如果有 await,并发请求可能在中间插进来。另一种是先预扣再判幂等:如果某段进程内判定依赖「单线程事件循环里两次 await 之间不会被别的协程打断」,查库的预扣就只能排在这段之前,代价是进门就登记,并发超限、已在运行、重复请求这几条被拒的路径都要各自回滚,把预扣、在跑位置和指纹还回去。注意「没有 await」只在单进程内成立,多副本下判定本身要靠数据库条件 UPDATE 或 Redis SET NX 这类原子操作,这时两者的先后可以按成本来排。
  • “任务在队列里等了很久没人领,怎么处理?” → 设入队等待上限,超时后先落取消标记再复查状态(防止刚好被领走),释放预扣和占位,终态写 failed 让用户可以重发。

易错:

  • ❌ “用了消费者组 + ACK 就是 exactly-once” → 只是 at-least-once,ack 前崩溃必然重投,幂等要业务层做。
  • ❌ “worker 被掐断的任务留在 PEL 等重跑最安全” → 重跑时用户可能早已离开,事件没人看,模型调用却真的再花一次;要按场景判断重跑还是给中断定论。
  • ❌ “多副本下用进程内 dict 记录谁在跑、做去重” → 每个副本各一份,同一会话能在两台机器上各起一个任务,必须放到 DB 条件更新或 Redis。