面试知识库

07 · 任务调度#

简历原话:任务调度:基于 Redis Stream 消费者组实现 at-least-once 任务队列;执行期间通过 lease key 定时续租,租约过期后由其他 worker 接管,任务版本号防止旧 worker 重复写入。按步保存执行状态,接管后从中断处继续,不重复扣费和加购。实测 kill -9 任一 worker,运行中的任务无丢失。

30 秒口述版#

ShoppingX 一轮购物任务要跑 4 次左右模型调用加一串工具,十几秒到几分钟,所以我把「收请求」和「跑 Agent」拆成 API 和 worker 两种进程,中间用 Redis Stream 消费者组做 at-least-once 队列。worker 领到消息就占一个 30 秒的租约键,每 10 秒用 Lua「值是自己才续期」,进程死了租约自然过期,空闲的 worker 抢到租约再 XCLAIM 接管。接管时发一个新的任务版本号,所有写入都带版本号做条件更新,被判死但其实还活着的旧 worker 写不进去。执行状态按步存进 MySQL,接管方从最后一个检查点接着跑,不重调已经算过的模型;扣费按 run_id 只结算一次,加购按「run_id + 动作 + 载荷指纹」唯一键只出一张确认卡。在 2 API + 2 worker 的多副本环境里反复 kill -9 任一 worker,运行中的任务无丢失。

背景与问题#

最早的形态是 API 进程收到 POST /api/task 后直接在本进程里起一个协程跑 AgentLoop。单机够用,但上线后碰到三类具体问题:

  1. 发布一次丢一批任务。每次滚动更新或 OOM 重启,进程里正在跑的任务直接没了,用户那条 WebSocket 停在「思考中」再也不动,刷新页面只能看到自己那句话,没有回答。
  2. 扩容只能整进程扩。收请求是轻活(鉴权、配额、入队,几毫秒),跑 Agent 是重活(一串 LLM 外呼 + 向量检索),两者资源曲线完全不同,绑在一起只能按重的那头买机器。
  3. 削不了峰。突发流量只能 429,没有地方让任务排一会儿再跑。

拆成队列后又冒出第二层问题,这才是这条 bullet 真正要解决的:

  • worker 崩了怎么办。Redis List 一 BRPOP 消息就从 Redis 消失,worker 跑到一半被 kill,这条任务无人知晓。
  • 怎么判断「崩了」。一开始用 Stream 的 idle 时间判死,阈值设 600 秒:短了会把正在跑的长任务误判成孤儿,长了用户要干等十分钟。跑着的任务 idle 本来就会一直涨,idle 根本不是「活着」的信号。
  • 判错了怎么办。worker 可能只是 GC 停顿、网络抖了一下、被 docker pause 冻住,租约过期后它又醒过来接着写,和接管方双写同一个会话。
  • 接管后从哪开始。整轮重跑意味着前面几步的模型调用再花一遍钱,如果中断前已经出过下单确认卡,重跑还可能让模型选出不同的商品、再出一张卡,用户就看到两张。

做法与取舍#

1. 队列:Redis Stream 消费者组,at-least-once#

  • 两条流 globex:intents(普通)和 globex:intents:large(长续聊),同一个消费者组 globex-workers。XREADGROUP 时普通流排在前面,同一批里短任务先拿到;长续聊不会被饿死,只是排在后面。
  • 消息领走后留在 PEL(pending 列表)里,直到 XACK 才算完。worker 崩了,消息还在 PEL,别人可以领回来。
  • 同一条消息投递 3 次仍失败就进死信流 globex:intents:dead(保留最近 1000 条)并 ack 掉,防止一条注定失败的消息反复重投、把 worker 并发度吃干净。
  • 背压:入队前用一段 Lua 把「未投递数 + 已准入未入队数」判定和登记做成一步,超过上限直接 429(这块属于限流那条 bullet,这里不展开)。
  • worker 内部用信号量限并发(默认 4)。XREADGROUP 的 COUNT 是每条流各自的上限,两条流一次最多返回 2×COUNT 条,只靠「少读几条」会跑出双倍并发。

为什么选 Redis Stream:项目已经有 Redis(事件回放、限流、控制面),Stream 自带消费者组、PEL、投递计数、XCLAIM,at-least-once 需要的东西都有。 否掉 Redis List:没有 PEL,弹出即消失,崩溃后无法重投。用 BRPOPLPUSH 自己维护一个处理中列表也能做,但投递计数、按消费者查挂起项都得自己写,等于重造一遍 Stream。 否掉 Kafka / RabbitMQ:任务量是每秒个位数到几十,Kafka 的分区和再平衡模型对这种「单条消息跑几分钟」的长任务不友好(一个慢消息卡住整个分区的位点提交),而且要多运维一个集群。RabbitMQ 的 ack 超时语义可以做,但同样是多一个组件,收益只是省掉几十行接管逻辑。

2. 判死:租约键 + 心跳续期,不看 idle#

  • 把「消息投递」和「任务存活」拆开:PEL 只管投递(谁领了、投了几次),独立的租约键 globex:lease:<流名>:<消息ID> 管存活,值是 worker 名(主机名 + pid,或 K8s pod 名),过期时间 30 秒。
  • 领到消息立刻占租约,不是等拿到信号量才占:在信号量上排队的消息同样挂在 PEL 里,没租约会被别人当孤儿领走。
  • 心跳每 10 秒跑一次 Lua:GET 出来的值是自己才 PEXPIRE,否则返回 0。租约是心跳的 3 倍,一次网络抖动没续上不会出事。
  • 接管方只在空闲时扫:XPENDING 带 IDLE 筛出领走超过 10 秒的候选(在跑的长任务也会被筛出来),逐条 SET NX 抢租约,抢到了说明原持有者没在续期,再 XCLAIM;抢不到就是还活着,跳过。
  • 任务结束(成功、失败、被取消)都用 Lua「值是自己才 DEL」释放租约,失败的任务能立刻被重投,不用干等 30 秒。

为什么续期要用 Lua:GET 和 PEXPIRE 分成两条命令,中间租约过期、被别人占走,续期就续成了别人的。和 Redisson 看门狗是同一个写法(它默认也是 30 秒租约、10 秒续一次)。 为什么先抢租约再 XCLAIM:两个接管方同时看到租约没了,只有 SET NX 成功的那个往下走;XCLAIM 再带 min-idle-time 挡一次(第一个领走后 idle 清零,第二个领不到)。 否掉 XAUTOCLAIM:它只看 idle,会把正在续租的长任务直接领走。 否掉「用 XCLAIM JUSTID 把 idle 清零当心跳」:这是上一版的做法。「查属主 → XCLAIM」是两条命令,中间被接管的话,心跳会把消息从接管方手里抢回来,两边都以为自己是持有者。

3. 任务版本号:挡住「以为自己还活着」的旧 worker#

租约只能保证「同一时刻最多一个人认为自己该跑」在大多数时候成立,保证不了全部。典型反例:worker A 卡了 40 秒(长 GC、宿主机被抢 CPU、容器被 pause),租约过期,worker B 接管;A 醒来后并不知道自己被判死,下一行代码就是写检查点、推事件。所以在租约之外再加一层版本号(也就是 fencing token):

  • MySQL 表 run_checkpoints 每个 run 一行,带 epoch 列。每次有 worker 领到这个任务(首次领取或接管),先执行 UPDATE run_checkpoints SET epoch = epoch + 1, owner = 自己 WHERE run_id = ?,读回的新值就是它这一任期的版本号。首次领取时行不存在,用 INSERT ... epoch = 1 建出来,唯一键冲突说明别人抢先了,改走接管分支。
  • 之后这个 worker 的每一次持久写入都带上自己的版本号做条件更新:写检查点是 WHERE run_id = ? AND epoch = 我的版本号;写 Redis 里的任务状态键、往事件流追加事件,走 Lua 先比对镜像键 globex:run:<task_id>:epoch,值相等才写。镜像键只由发号成功的新持有者用「只增不减」的 Lua 设置。
  • 影响 0 行 / Lua 返回 0,就说明自己已经被接管。旧 worker 此时立即停止这条任务:取消本地协程,不写终态,也不 XACK,只打一条「已被接管,放弃」的日志。

为什么旧 worker 不能 XACK:XACK 不校验属主。如果旧 worker 醒来后把消息 ack 掉,这条消息就从 PEL 里消失了;万一接管方随后也崩了,就没人能再领回来,任务真的丢了。 为什么版本号放 MySQL 不放 Redis INCR:检查点本来就在 MySQL,版本号和检查点在同一行,条件更新一条语句同时完成「校验身份 + 写状态」,没有两个存储之间的时间窗。Redis 主从切换可能丢最近一秒的写,INCR 出来的号有回退风险;MySQL 这边是事务提交后才算数。 否掉「只靠租约不加版本号」:租约的正确性建立在「进程停顿不会超过租约时长」这个假设上,这个假设没法保证。Martin Kleppmann 批评 Redlock 用的就是这个例子:锁过期后客户端还以为自己持锁,只有存储端校验单调递增的 token 才挡得住。 否掉「续期失败就立刻 cancel」:只做这一步不够。心跳和主任务是两个协程,心跳发现失去租约时,主任务可能正处在一次写入的半路;而且卡住的恰恰可能是整个进程,心跳协程自己也没机会跑。存储端校验不依赖旧 worker 的自觉。

4. 按步保存执行状态:从中断处继续#

「一步」怎么切:AgentLoop 天然的边界是「一次模型调用 + 它发出的这批工具调用全部拿到结果」。普通轮大约 4 步(planner → item_search → 框架自动比价精挑 → shopping_summary)。每一步存两次:

  • 检查点 A(决策后):模型返回、工具调用列表确定之后,工具执行之前。存下模型这次的输出(含每个工具调用的 id 和参数)和这次调用的 token 用量。
  • 检查点 B(结果后):这批工具全部有结果之后。存下工具结果,步号 +1。

存什么:一个检查点就是「重建这一刻所需的全部状态」,压缩后几十 KB:

  1. AgentScope 的 Agent 状态对象(消息列表、当前迭代数、上下文压缩摘要)。
  2. 控制面的轮内状态:候选商品登记表(下单工具只收登记过的商品 id)、自动比价的武装位、检索预算的三本账、循环检测计数。这些不存的话,接管后预算被重置为满额,模型可以借一次崩溃多搜一轮;登记表丢了,后面下单会因为「商品不在登记表里」被拦截。
  3. 已累计的模型用量(输入 / 输出 token、折算金额)。
  4. 事件序号:最后一条已推给前端的事件 ID。

接管后怎么续,看最后一个检查点是哪一种:

  • 最后是 B:上一步完整结束,直接从下一步的模型调用开始。
  • 最后是 A:模型已经决定调哪些工具,但工具没跑完。不重新调模型,按检查点里记下的调用列表重新执行这批工具:只读工具(检索、比价、运费)直接重跑;写工具走幂等键(见下一节),重跑拿到的是第一次的结果。
  • 一个检查点都没有:从头跑,和第一次一样。

事件怎么接上:接管方先推一条 task_resumed 事件,带上检查点里的事件序号,前端把这个序号之后、属于中断那一步的半截内容收起。重跑的工具沿用检查点 A 里的工具调用 id,前端按 id 合并工具卡片,不会出现两张「正在检索 Amazon」。

为什么存 MySQL 不存 Redis:检查点是「任务不丢」这件事的最后一道凭据,Redis 按 AOF everysec 配置最多丢 1 秒,恰好是崩溃现场那一秒。写入频率是每任务每步 2 次,一轮 8 次左右,一次几毫秒,对十几秒的任务无感。 为什么切两个检查点而不是一个:只在 B 存,崩在工具执行中就得重调模型;模型这次可能选出不同的商品,下单卡的载荷指纹就变了,幂等键挡不住,用户会看到两张卡。在 A 先把「决定」存下来,重放的是同一份决定。 否掉「只在整轮结束时存会话文件」:原来会话状态只在整轮结束时写一次,而且写在本机磁盘,接管方所在的机器根本读不到。 否掉照搬 LangGraph checkpointer:思路一致(按步存、按 thread 恢复),但项目主循环是 AgentScope,状态对象和步边界都不同;需要的只是「一张表 + 两个挂点」,引一整套图执行框架不划算。

5. 不重复扣费、不重复加购#

at-least-once 意味着「同一条任务可能被执行不止一次」,所以所有有副作用的地方都要幂等。本 bullet 关心的是两处:

扣费:

  • 入队前按档位预扣一笔(普通 20、长续聊 60 credits),写一行 run_holds,主键就是 run_id,而 run_id 恒等于 task_id,重投、接管都不会变。
  • 用量不在每次模型调用后扣,而是随检查点累计。接管方从检查点读出已累计的用量接着加,崩溃那一步里「调了模型但还没存进检查点 A」的那次调用,不计入用户的账(平台自己承担,一次崩溃最多一次调用)。
  • 任务结束时结算一次:UPDATE run_holds SET state = 'settled', credits_charged = ? WHERE run_id = ? AND state IN ('queued', 'running'),和记账写在同一个事务里。影响 0 行说明已经结算过,直接跳过。
  • 预扣行 15 分钟过期。kill -9 之后如果一直没人接管(比如所有 worker 都挂了),过期后额度自动释放,不会一直占着用户的并发名额。

加购(下单确认卡):

  • 下单工具不直接下单,只出一张确认卡,用户点了确认才通过 HTTP 真正落订单。
  • 确认卡表有一个唯一键 request_key = run_id : 动作 : 载荷指纹。载荷指纹是把商品按 id 排序后再算哈希,模型把同样两件商品换个顺序报上来,指纹不变。
  • 接管后重放检查点 A 里的同一次下单调用,参数完全相同,撞上唯一键,直接返回第一次那张卡。
  • 就算崩在「用户已点确认、订单还没写完」这一刻,落订单那一步走的是确认记录自己的操作 ID 做幂等键,重复提交只会得到同一个订单号。

否掉「每次模型调用后立即扣费」:扣费和检查点不在一个事务里,崩在两者之间,要么扣了没记检查点(接管后再扣一次),要么记了没扣。随检查点累计、结束时一次结算,只有一个写入点。 否掉「用 tool_call_id 做加购幂等键」:框架的工具中间件拿不到模型给的调用 id;而且整步重调模型时 id 会变,只有「run + 动作 + 内容」这个组合在重放时稳定。

6. 重投时的去重和另外两条收尾路径#

  • 写完终态、还没 XACK 就被 kill:消息会被接管方再领一次。领到手先查 Redis 里的任务状态,已经是 done / cancelled / interrupted 就直接 ack 跳过,不把已经答完的一轮再跑一遍。failed 不在跳过名单里,失败本来就要留 PEL 重投。
  • 优雅关停(SIGTERM):先停领新任务,再等在飞任务最多 330 秒(大于单轮超时 300 秒)自然跑完;超时的少数长尾按 interrupted 收尾、推事件让用户重发、ack 掉,不交给接管。理由是滚动发布时多个副本同时退出,一波接管会全部压到剩下的副本上,而这些任务本来就能在发布窗口里跑完。
  • kill -9 / 宿主机宕机:没有任何收尾代码能跑,这是接管续跑唯一覆盖的场景,也是简历里那个实测的对象。

流程图 / 架构图#

图 1:整体架构#

flowchart LR
    U["用户 / 前端"] -->|POST 任务| API["API 副本 x2"]
    API -->|预扣 credits| DB[("MySQL<br/>run_holds / run_checkpoints / 确认卡表")]
    API -->|XADD| S["Redis Stream<br/>globex:intents / large"](#)
    S -->|XREADGROUP 消费者组| W1["worker 1"]
    S -->|XREADGROUP 消费者组| W2["worker 2"]
    W1 -->|占租约 + 每 10 秒续期| L[("Redis 租约键<br/>globex:lease:流:消息ID")]
    W2 -->|空闲时扫 XPENDING<br/>SET NX 抢过期租约| L
    W1 -->|每步两次检查点<br/>带版本号条件写| DB
    W2 -->|接管: 发新版本号<br/>读检查点续跑| DB
    W1 -->|事件| E["Redis 事件流<br/>带版本号校验"](#)
    W2 -->|事件| E
    E -->|WebSocket 推送| U
    S -.->|投递 3 次仍失败| D["死信流"](#)

图 2:kill -9 后的接管时序#

sequenceDiagram
    participant R as Redis Stream + 租约
    participant A as worker A
    participant B as worker B
    participant M as MySQL 检查点
    A->>R: XREADGROUP 领到消息 m
    A->>R: 占 m 的租约, 值为 A, 30 秒
    A->>M: 发号 epoch = 1
    loop 每一步
        A->>M: 检查点 A 决策后, WHERE epoch = 1
        A->>A: 执行工具
        A->>M: 检查点 B 结果后, WHERE epoch = 1
        A->>R: 每 10 秒 Lua 续期
    end
    Note over A: 第 3 步工具执行中被 kill -9
    Note over R: 最多 30 秒后租约过期
    B->>R: 空闲, XPENDING 筛出 idle 超 10 秒的 m
    B->>R: SET NX 抢 m 的租约成功
    B->>R: XCLAIM m, 投递次数 +1
    B->>M: 发号 epoch = 2, 读最后检查点
    M-->>B: 第 3 步检查点 A, 工具调用清单
    B->>B: 按清单重放工具, 写工具命中幂等键
    B->>R: 推 task_resumed 事件
    B->>M: 继续第 3 步检查点 B, WHERE epoch = 2
    B->>M: 结算 run_holds, 只生效一次
    B->>R: 写终态 done, XACK m, 释放租约

图 3:旧 worker 醒来被版本号挡住#

sequenceDiagram
    participant A as worker A 被冻住 40 秒
    participant B as worker B 接管方
    participant M as MySQL 检查点
    A->>M: 写检查点 WHERE epoch = 1, 成功
    Note over A: 长时间停顿, 租约过期
    B->>M: 发号 epoch = 2
    B->>M: 写检查点 WHERE epoch = 2, 成功
    Note over A: 醒来, 不知道自己被接管
    A->>M: 写检查点 WHERE epoch = 1
    M-->>A: 影响 0 行
    A->>A: 取消本地协程, 不写终态, 不 XACK

图 4:接管后从哪一步开始#

flowchart TD
    S["接管方读最后一个检查点"] --> Q{"检查点类型"}
    Q -->|没有检查点| N["从头跑本轮"]
    Q -->|B 结果后| NB["从下一步的模型调用开始"]
    Q -->|A 决策后| NA["不重调模型<br/>按记下的工具清单重放"]
    NA --> T{"工具类型"}
    T -->|只读: 检索 比价 运费| RR["直接重跑"]
    T -->|写: 下单确认卡| WR["按 run_id + 动作 + 指纹<br/>命中唯一键, 返回原卡"]
    RR --> C["写检查点 B, 继续循环"]
    WR --> C
    NB --> C
    N --> C

数字怎么来的#

简历上这条只有一个结论型指标:kill -9 任一 worker,运行中的任务无丢失。它的口径如下。

环境:一套独立的多副本验收环境,docker compose 起 2 个 API 副本 + 2 个 worker 副本 + MySQL 8 + Redis 7,和线上是同一个镜像、同一条构建路径。租约 30 秒、心跳 10 秒、接管扫描 idle 阈值 10 秒、单 worker 并发 4,都是线上默认值。

「任务无丢失」的定义,每个被测任务同时满足四条才算通过:

  1. 有终态:GET /api/task/{id} 最终拿到 done,且回答内容完整(有清单、有理由),不是空串。
  2. 扣费一次:run_holds 里这个 run 恰好一行 settled,扣掉的 credits 和同一条 query 不 kill 的对照跑相比,差值不超过「一次模型调用」的量。
  3. 加购一次:带下单意图的 query,确认卡表里按 request_key 恰好一张卡。
  4. 事件不乱:前端按事件 ID 去重后,工具卡片没有重复,task_resumed 之前的半截步骤被正确收起。

怎么杀:

  • 用的是 docker kill -s KILL,等价 kill -9,进程没有任何机会跑收尾代码。
  • 「任一」是说两个 worker 轮流杀,每轮杀的是当时手上在跑任务最多的那个。
  • 杀的时机按步随机:模型调用进行中、工具执行中、两个检查点之间、写完终态还没 XACK,四种都要覆盖到。杀在哪一刻由压测脚本读事件流判断(看到第 N 步的 tool_start 就动手)。
  • 被杀的 worker 过几秒再拉起来,模拟 K8s 重启 Pod。

跑法和样本量:

  • 桩模型:把模型调用换成一个可控的桩,每步固定耗时、固定 token 用量、第 3 步固定发一次下单调用。这样「扣费一致」可以精确比对,杀的时机也能精确控制。跑了 30 轮,每轮同时提交 6 个任务,共 180 个任务,全部满足四条。
  • 真模型:用真实 LLM 跑 20 条带下单意图的购物 query,每条中途杀一次,20 条全部 done,每条一张确认卡。真模型下扣费只核对「一行 settled」,不比对金额(真模型每次用量本来就有波动)。
  • 僵尸 worker 专项:用 docker pause 把持有任务的 worker 冻住 45 秒(超过 30 秒租约)再恢复,验证版本号挡写。10 次全部是:接管方跑完,旧 worker 醒来后第一次写检查点影响 0 行、自行放弃、没有 XACK。

附带测到的接管延迟(不在简历上,被问到可以说):从 kill 到接管方推出 task_resumed,主要由 30 秒租约决定,实测大约 31~44 秒;多出来的部分是接管方要等到自己空闲时才扫 PEL,外加 XREADGROUP 的 2 秒阻塞。

对比基线:接管机制之前是纯 idle 判死(600 秒)+ 整轮重跑,同样的 kill 测试任务最终也能跑完,但要等 10 分钟以上,且带下单意图的 query 有一部分出了两张确认卡(重跑时模型选了不同商品)。

追问 Q&A#

Q1:为什么说 at-least-once,不直接做 exactly-once?

分布式系统里「消息只投递一次」做不到:worker 处理完、ack 之前崩了,队列无法区分「没处理」和「处理完没来得及 ack」,只能再投一次。能做的是「至少投一次 + 消费端幂等」,效果上等于恰好一次。我这里的幂等分三层:终态去重(已经 done 的直接 ack 跳过)、结算按 run_id 条件更新、加购按 request_key 唯一键。

Q2:租约 30 秒、心跳 10 秒是怎么定的?

两个约束。租约至少是心跳的 3 倍,一次网络抖动或一次慢的 Redis 往返没续上,不至于被误判死亡;消费循环启动时会检查这个比例,不满足就打告警。租约又不能太长,它直接决定 kill -9 后用户要等多久才恢复。30 / 10 和 Redisson 看门狗的默认值一样。注意租约不需要长于任务耗时,任务跑多久都靠心跳续着,这是和「纯 idle 判死」最大的区别。

Q3:有租约了为什么还要版本号?租约不就是分布式锁吗?

租约只在「持有者不会停顿超过租约时长」的前提下成立。GC、宿主机抢占、容器被冻住都可能打破这个前提:旧 worker 醒来时不知道租约已经过期,照样写。版本号是让存储端去判断「你还是不是当前持有者」,不依赖旧 worker 自觉。这就是 fencing token 的思路,Kleppmann 批评 Redlock 时举的正是这个例子。

Q4:Redis 主从切换,租约键丢了怎么办?

新主上没有这个租约键,另一个 worker 扫到后会以为原持有者死了、接管,这时两个 worker 同时在跑同一个任务。版本号在 MySQL,不受 Redis 切换影响:接管方发号拿到 epoch + 1,原持有者下一次写检查点影响 0 行,自己放弃。代价是原持有者那一步的模型调用白做了,但不会有双写,也不会重复扣费和加购。

Q5:接管只在 worker 空闲时才扫 PEL,高峰期所有 worker 都满载,孤儿任务会不会一直没人管?

会有延迟。满载时每个 worker 都在信号量上等,新消息都读不进来,确实没空扫。两点缓解:一是每个 worker 有空位时,XREADGROUP 读不到新消息才转去扫 PEL,高峰过去立刻会扫到;二是被接管的任务是「已经跑了一半的用户任务」,优先级本来就应该高于新任务,所以扫描还有一个定时触发:距离上次扫描超过 60 秒就强制扫一次,扫到的孤儿插到本地待执行队列最前面。

Q6:崩溃那一步已经调过的模型,钱谁出?这还算「不重复扣费」吗?

对用户来说算。用户的账只按检查点里累计的用量结算,崩在「模型调用中」或「模型返回了但还没写检查点 A」这一小段,那次调用没有进检查点,接管方会重新调一次,用户只为重调的那次付钱。第一次的费用平台自己承担,一次崩溃最多损失一次模型调用。我在检查点 A 就把用量记进去,就是为了把这个窗口压到最小。

Q7:接管后重放只读工具,结果可能和第一次不一样(比如价格刚变了),会不会前后矛盾?

不会矛盾,因为第一次的结果没有被任何人用过:检查点 A 之后、B 之前,工具结果还没回到模型那里,模型没据此做过任何决策;前端虽然可能已经显示了部分结果,但接管方推的 task_resumed 会把半截那一步收起、用重放结果替换。只有写工具需要严格一致,所以写工具走幂等键返回第一次的结果。

Q8:任务正停在 ask_user 等用户回复时 worker 挂了,怎么办?

ask_user 本身就是一个工具调用,发出去之前已经写了检查点 A。接管方重放它时,先按本轮的 turn_id 去 Redis 查有没有已经到达的回复:用户在 worker 挂掉期间就回复了的,回复会被控制面暂存,重放直接拿到,不再问第二遍;还没回复的,重新挂起等待,前端那张问题卡片按工具调用 id 合并,不会出现两张一样的问题。

Q9:一条任务每次跑都把 worker 搞崩(毒消息)怎么办?

每次接管用的是不带 JUSTID 的 XCLAIM,投递计数会 +1。计数到 3 仍然失败就写进死信流并 ack 掉,任务状态标 failed,推错误事件给前端。死信只保留最近 1000 条,是给人排查用的,不自动重放。不设死信的话,一条毒消息会被无限接管,每次都带走一个 worker。

Q10:为什么不用 Temporal 这类工作流引擎?它天然支持持久化执行和重放。

考虑过,否掉有两个原因。一是 Temporal 的重放要求工作流代码是确定性的,而 AgentLoop 的核心就是一次非确定的模型调用,要把每次模型调用都包成 Activity,等于把 AgentScope 的主循环拆开重写。二是多一个有状态集群(Temporal Server + 它自己的存储)要运维,而我需要的只是「消息不丢 + 按步存状态 + 防双写」,Redis Stream + 一张检查点表 + 版本号就够了。思路上是借鉴它的:检查点 A 存「决定」、重放时不重新决定,就是 Temporal 里「Activity 结果进历史、重放时直接读历史」的简化版。

Q11:检查点每步写两次 MySQL,会不会成为瓶颈?

不会。一轮大约 4 步、8 次写,每次压缩后几十 KB、按主键更新,几毫秒;而一步本身是秒级(一次模型调用)。按线上并发算,检查点写入是每秒几十次的量级,MySQL 完全吃得下。真到瓶颈,第一步是只保留最近一个检查点(现在就是覆盖写,不是追加),第二步才是换存储。

相关八股#

1. Redis Stream 消费者组的核心概念有哪些?

  • XADD 追加消息,消息 ID 是「毫秒时间戳-序号」,单调递增。
  • XGROUP CREATE 建消费者组,起始 ID 用 0 从头消费、用 $ 只消费之后的新消息。
  • XREADGROUP ... > 读从未投递给本组的新消息;同一组内一条消息只投给一个消费者,组与组之间各消费一遍。
  • 投出去的消息进 PEL(Pending Entries List),记录属主、idle 时间、投递次数,XACK 后才移出。
  • XPENDING 查挂起项,XCLAIM / XAUTOCLAIM 把 idle 超阈值的消息转给别的消费者,投递次数 +1(带 JUSTID 不加)。

本项目:两条流共用一个组;判死不用 idle 而用独立租约,所以接管用 XPENDING + SET NX + XCLAIM 组合,不用 XAUTOCLAIM。

2. 消息投递语义:at-most-once、at-least-once、exactly-once 的区别?

  • at-most-once:先 ack 再处理,崩了就丢。
  • at-least-once:处理完再 ack,崩在两者之间会重投,可能重复。
  • exactly-once:端到端只靠队列做不到,实际做法是 at-least-once + 消费端幂等(唯一键、条件更新、去重表),或者像 Kafka 那样把「写结果」和「提交位点」放进同一个事务。

本项目:at-least-once 投递,三层幂等(终态去重、run_id 结算、request_key 唯一键)。

3. Redis 分布式锁怎么实现?锁过期了业务没跑完怎么办?

  • 加锁:SET key 唯一值 NX PX 过期时间,一条命令保证原子。
  • 解锁:Lua 脚本先比对值是自己的再 DEL,避免删掉别人的锁。
  • 续期:看门狗线程定期检查,值是自己的就 PEXPIRE,Redisson 默认 30 秒锁、每 10 秒续一次。
  • 主从切换时锁可能丢,Redlock 试图用多数派解决,但仍然挡不住客户端长时间停顿。

本项目:租约键就是一把带看门狗的锁,加锁、续期、释放三段全是 Lua 比值后再动手。

4. 什么是 fencing token?为什么说单靠分布式锁不安全?

  • 锁服务每次发锁时给一个单调递增的号,客户端写存储时带上这个号,存储端拒绝比自己见过的号小的写入。
  • 原因:锁有过期时间,持锁客户端可能因为 GC、网络、调度停顿超过过期时间,醒来后仍以为自己持锁;只有存储端校验才可靠。
  • 要求存储端能做「比较后写入」,数据库的条件更新、ZooKeeper 的 zxid、etcd 的 revision 都可以当 token。

本项目:任务版本号就是 fencing token,由 MySQL 检查点行发号,写检查点和写 Redis 状态都做比较。

5. 乐观锁和悲观锁的区别?MySQL 里怎么实现乐观锁?

  • 悲观锁:先锁再改,SELECT ... FOR UPDATE,适合冲突多的场景,代价是持锁期间别人全阻塞。
  • 乐观锁:不加锁,更新时带条件 UPDATE ... SET version = version + 1 WHERE id = ? AND version = ?,影响 0 行说明被别人改过,重试或放弃。
  • 条件更新的原子性由 InnoDB 行锁保证:UPDATE 会对命中行加排他锁,判断条件和写入在同一把锁内完成。

本项目:写检查点是乐观锁写法,结算也是(WHERE state IN ('queued','running')),影响 0 行即「已被接管」或「已结算过」。

6. Redis 的 Lua 脚本为什么是原子的?有什么注意点?

  • Redis 执行命令是单线程的,一段 EVAL 脚本执行期间不会穿插其他客户端的命令。
  • 脚本不回滚:执行到一半报错,前面已经执行的写不会撤销。
  • 脚本要短:执行期间整个 Redis 被占住,长脚本会阻塞所有请求;超过 lua-time-limit 只会告警,不会中断。
  • 集群模式下脚本用到的键必须在同一个槽,可以用 hash tag {...} 强制同槽。

本项目:租约的领取 / 续期 / 释放、准入闸、版本号镜像的只增写入都是短 Lua。

7. Redis 持久化 RDB 和 AOF 的区别?对队列可靠性有什么影响?

  • RDB:定时全量快照,恢复快,但两次快照之间的写会丢。
  • AOF:追加写命令日志,appendfsync 可选 always / everysec / no,everysec 最多丢 1 秒。
  • 混合持久化:AOF 重写时前半段写 RDB 格式,兼顾恢复速度和数据完整。
  • 主从复制是异步的,主挂了、从还没收到的写会丢。

本项目:Stream 开 AOF everysec,能容忍极小概率丢最后 1 秒的入队;但「任务跑到哪了」这种丢了就要重跑的状态放 MySQL,不放 Redis。

8. Redis Stream 和 Kafka 做任务队列怎么选?

  • Kafka:分区内有序、按位点提交,吞吐极高、可长期保留、可回放;但一个分区同一时刻只给组内一个消费者,慢消息会卡住后面整个分区的位点提交,不适合单条要处理几分钟的任务。
  • Redis Stream:逐条 ack、PEL 可按条接管,适合「单条慢、总量不大」的任务队列;缺点是数据在内存里、容量受限,持久化和复制的可靠性不如 Kafka。
  • RabbitMQ:逐条 ack、有死信交换机,也适合任务队列,多一个组件要运维。

本项目:每秒几十条、单条几十秒到几分钟、已经有 Redis,选 Stream。

9. 什么是检查点 / 断点续跑?常见实现有哪些?

  • 思路:把「重建当前进度所需的全部状态」定期持久化,失败后从最后一个检查点恢复,而不是从头来。
  • Flink 的 checkpoint 按 barrier 对齐做全局快照;LangGraph 的 checkpointer 每个节点执行后按 thread 存一份状态;Temporal 把每个 Activity 的结果写进事件历史,恢复时重放代码、直接读历史结果。
  • 关键点:步边界怎么切、有副作用的步骤怎么保证重放不重复(幂等或只重放记录下来的决定)。

本项目:步边界是「一次模型调用 + 这批工具全部有结果」,每步存决策后、结果后两个检查点,重放记录下来的工具调用而不是重新调模型。