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 过的条目 |
import asyncio
import redis.asyncio as redis
STREAM, GROUP = "agent:tasks", "workers"
running: set[asyncio.Task] = set() # 本 worker 在跑的任务,SIGTERM 时等它们
slots = asyncio.Semaphore(MAX_CONCURRENCY) # 单 worker 并发上限
async def worker_loop(r: redis.Redis, consumer: str, stop: asyncio.Event):
while not stop.is_set():
await slots.acquire()
resp = await r.xreadgroup(GROUP, consumer, {STREAM: ">"}, count=1, block=5000)
if not resp:
slots.release()
continue
for _, messages in resp:
for msg_id, fields in messages:
t = asyncio.create_task(handle(r, msg_id, fields))
running.add(t)
t.add_done_callback(lambda t: (running.discard(t), slots.release()))
async def handle(r, msg_id, fields):
task_id = fields[b"task_id"].decode()
if not await claim_run(task_id): # 业务幂等:只有 queued 能被领成 running
await r.xack(STREAM, GROUP, msg_id) # 重复投递,直接确认
return
try:
await run_agent(task_id) # 正常结束时内部写 succeeded / failed 终态
except asyncio.CancelledError:
await finalize_stopped(task_id) # 用户取消写 cancelled、排空超时写 interrupted,按取消标记区分
except Exception:
log.exception("run %s failed", task_id) # 单个任务出错不打断领取循环;failed 终态由 run_agent 写
await r.xack(STREAM, GROUP, msg_id) # 终态写完才 ack
async def reclaim_loop(r, consumer):
# 接管崩溃 worker(kill -9 / OOM)留下的消息。这里演示「不重跑、写中断终态后 ack」这种取舍;
# 另一种是把领回的消息交给 handle 重跑,需要配合投递次数上限和死信(见追问)
# min_idle 相当于租约,必须大于单任务最长执行时间,否则会把正在跑的任务当成崩溃
while True:
_, msgs, *_ = await r.xautoclaim(STREAM, GROUP, consumer,
min_idle_time=MAX_RUN_MS + 60_000, count=10)
for msg_id, fields in msgs:
# 发中断事件 → 释放占位和预扣 → UPDATE runs SET state='interrupted'
# WHERE id=? AND state IN ('queued','running')
await finalize_interrupted(fields[b"task_id"].decode())
await r.xack(STREAM, GROUP, msg_id)
await asyncio.sleep(30)python长任务专有的坑:PEL 里的「空闲时间」从投递那一刻算起,任务跑着但没 ack 也在累加。min-idle 如果小于任务的最长执行时间,另一个 worker 会把正在跑的任务当成崩溃接管:按上面的写法会被误判成中断,如果接管后重跑,就是同一个任务跑两遍。两种解法:min-idle 设成大于单轮超时;或者执行中的 worker 定期对自己手里的消息 XCLAIM ... JUSTID 给自己续约(认领会重置空闲时间)。
3. 幂等:队列保证不了 exactly-once#
重复投递来源:worker 执行完、ack 前崩溃;XAUTOCLAIM 误领;客户端超时重发。通用手段见 消息重复与幂等方案,Agent 场景的三层:
- 入口去重:请求指纹
SET dedup:{hash} 1 NX EX <窗口秒数>,查和登记是一步原子操作。Redis 不可达时有两种取舍:直接报错(503),代价是可用性;或者退回进程内字典,代价是多副本下重复请求会漏过去。涉及扣费和写操作时一般选前者。 - 任务状态机条件更新:领任务时
UPDATE runs SET state='running', worker=? WHERE id=? AND state='queued',影响行数为 1 才执行;判定和占位在同一条语句,数据库保证并发两条只有一条成功。已经是running的任务被重投时不会再跑一遍,交给崩溃接管的策略处理(见第 5 节)。 - 副作用幂等:下单、发邮件这类工具调用带
operation_id,下游按它去重;结算按run_id幂等(见第 6 节)。
4. 和 Celery / Kafka / RabbitMQ 的取舍#
| 方案 | 适合 | 长任务要注意的点 |
|---|---|---|
| Redis Stream | 已有 Redis、量不大、想少一个中间件 | 数据在内存要 MAXLEN;min-idle 与任务时长的关系要自己管;没有内置延迟队列和死信 |
| Celery | Python 生态、要现成的重试/定时/结果后端 | 默认执行前 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)pythonGRACE_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 (发布/排空超时被掐)
└──超时无人领取──▶ failedplaintext- 每次状态迁移都用条件更新
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。