← 返回笔记

learning note

Redis 状态机设计思路

把分布式 Single-flight 从“有锁或没锁”升级为一次请求协调实体的完整生命周期管理。

为什么需要状态机?

单机版 Single-Flight 比较简单:同一个 key 来了之后,当前 JVM 里只要保存一个 CompletableFuture,后来的线程直接等这个 future 完成即可。

但分布式场景不一样。请求可能被负载均衡打到不同机器,这时系统要解决的不只是“当前线程要不要执行”,而是下面这些更复杂的问题:

  • 哪个节点抢到了真正执行权

  • 哪个节点只是 follower,应该等待结果

  • owner 是刚创建的,还是对旧任务的接管者

  • owner 现在是否还活着

  • 请求最后是成功、失败、取消,还是已经过期

  • 失败后能不能重试,能不能由别的节点接管

  • follower 现在应该等待、复用结果,还是直接失败返回

这些问题本质上都不是一个布尔值能表达的,所以必须把一次 flight 抽象成一个 有生命周期的状态机

一句话说:

分布式 Single-Flight 不是“有锁/没锁”这么简单,而是“一个请求协调实体从创建、运行、完成到失效的完整生命周期管理”。

需要注意的是:状态机本身没什么难的。他就是人为规定一个流程的多个状态,然后定义状态之间的流转规则而已

状态机落在哪里

这套状态机不是只存在于 Java 内存里,而是持久化在 Redis 里,主要就放在这个Key中:

ai:flight:meta:{requestKey}

状态机的主体保存在 meta hash 中,常见字段包括:

  • status:当前 flight 处于什么状态,比如 PENDING、RUNNING、SUCCEEDED、FAILED。它是整个状态机的核心字段。

  • stage:这次 AI 请求属于哪个业务阶段,比如 interview-extraction、interview-demeanor、interview-evaluation。它表示这是“哪一类业务”的 flight。

  • ownerId:当前真正持有执行权的节点是谁,一般是某台机器或某个进程实例的标识。

  • ownerToken:当前 owner 的版本号,也可以理解成“这一任 owner 的执照编号”。它主要用来防止旧 owner 恢复后把过期结果写回来。

  • requestKey:这次 flight 对应的唯一业务指纹,也就是 single-flight 真正拿来判定“是不是同一个请求”的 key。

  • sessionId:这次请求所属的面试会话 ID,用来把 flight 和具体业务会话关联起来,便于隔离和排查。

  • createdAt:这条 flight 元数据第一次被创建的时间。

  • updatedAt:这条 flight 最近一次被更新的时间,比如状态变化、heartbeat 续租、失败落库时都会更新它。

  • heartbeatAt:最近一次心跳时间。follower 会根据它判断 owner 还活着,还是已经失活可以接管。

  • expireAt:这条 flight 逻辑上的过期时间。它表示这条记录预计什么时候失效,不该再继续被当成活跃执行看待。

  • retryable:如果这次执行失败了,后续请求能不能重试、能不能接管。true 表示可以,false 表示失败就是最终结论。

  • errorType:失败的大类,比如超时、过载、上游错误、参数错误。它是失败原因的分类标签。

  • errorCode:更具体的失败码,用来表达更明确的错误结论,方便日志、排查和返回前端。

  • resultRef:结果引用,指向真正结果存储的位置。当前实现里通常就是 Redis 里的 result key,用于让 follower 去读最终结果。

也就是说,当前状态机不是靠内存变量维持,而是靠 Redis 里的元数据实体维持。这样无论请求落到哪台机器,只要大家读写的是同一个 requestKey 对应的 meta,就能共享同一套状态。

当前状态集合

当前代码里定义的状态有六个:

  • PENDING

  • RUNNING

  • SUCCEEDED

  • FAILED

  • CANCELLED

  • EXPIRED

其中前四个是当前主链路里真正会用到的核心状态;后两个是为了把状态机定义补完整,当前实现中更多属于 保留态/扩展态

每个状态分别表示什么

PENDING

PENDING 表示:

这个请求协调实体已经被创建出来了,owner 也已经被选出来了,但真正的业务执行还没有正式进入运行态。

在当前实现里,以下情况会把状态置为 PENDING

  1. 这个 key 第一次出现,acquireOrJoin(...) 抢占成功

  2. 之前是 FAILED 且允许重试,新的 owner 接管成功

  3. 之前 owner 心跳超时,被新 owner takeover

为什么需要 PENDING,而不是一创建就直接 RUNNING

因为“抢到执行权”和“真正开始执行业务”是两个不同的瞬间:

  • 抢占阶段只是声明“我来负责这个 flight”

  • 运行阶段才表示“我已经进入真实执行,并开始 heartbeat”

所以 `PENDING` 的存在,本质上是在状态机里显式区分 占坑成功正式开跑 这两个动作。

RUNNING

RUNNING 表示:

owner 已经正式开始执行 supplier,flight 处于活跃运行中。

一旦进入 RUNNING

  • owner 会周期性 heartbeat

  • follower 会把它视为“当前已有活跃执行者”

  • 旧 owner 如果后来恢复,也不能再写回结果,因为会被 ownerToken 拦住

RUNNING 是整套状态机里最关键的活跃态,因为它对应的是“真实模型调用正在发生”。

SUCCEEDED

SUCCEEDED 表示:

owner 已经把结果写入 result 存储,并完成成功收尾。

进入这个状态之后:

  • follower 不需要再等待

  • 后续同 key 请求可以直接走回放

  • L1 本地缓存也可以被填充

这个状态是典型的终态,主要服务于 结果复用

FAILED

FAILED 表示:

owner 执行失败,且已经把失败结论写入状态机。

这个状态不是简单地表示“出错了”,它还会配合两个字段一起看:

  • errorType

  • retryable

也就是说,FAILED 还会继续分成两类:

- 可重试失败:比如超时、过载、上游不可用

- 不可重试失败:比如参数错误、某些明确业务异常

如果是可重试失败,后续请求可以把它从 FAILED 重新推进到 PENDING,由新 owner 接管;如果是不可重试失败,后续请求会直接读取失败结论,而不是继续重复打模型。

CANCELLED(预留)

CANCELLED 表示:

这个 flight 被主动取消,不应该再继续执行。

不过要特别说明:

这个状态目前更多是为未来扩展预留的,比如:

  • 人工中止任务

  • 上层流程状态机主动取消

  • 用户离场后终止某些异步 flight

当前实现里,Lua 脚本已经认识 CANCELLED,但主链路暂时没有写入它的完整流程。

EXPIRED(预留)

EXPIRED 表示:

这个 flight 因超时或生命周期结束而失效,不应再被继续复用。

同样要强调一点:

当前代码里也没有完整的显式状态流把 meta 写成 `EXPIRED`。

当前实际的“过期”主要是通过 Redis TTL 删除 meta/result key 来完成的。也就是说,在现阶段的真实运行里,一个 flight 过期后,更多表现为:

  • meta 被 Redis 删除

  • 下次再看这个 key,像是一个全新的请求

所以 EXPIRED 在当前实现里也是一个预留态,比起“正在被主链路广泛使用的状态”,它更像“状态机设计上已经考虑到,但尚未完全落地的终态”。

当前状态机的主路径

如果只看当前代码里真正落地的主链路,最核心的状态迁移其实非常清晰:

PENDING -> RUNNING -> SUCCEEDED
PENDING -> RUNNING -> FAILED
FAILED(retryable=true) -> PENDING -> RUNNING -> SUCCEEDED/FAILED

也就是说,真正高频发生的是这几件事:

  1. 抢占成功,生成 PENDING

  2. owner 正式开跑,进入 RUNNING

  3. 执行完成后进入 SUCCEEDEDFAILED

  4. 如果失败可重试,再次接管,又重新回到 PENDING

一个完整例子

假设同一个 requestKey 的请求同时打到两台机器 A 和 B。

第一步:A 抢占成功

  • Redis 里原来没有 meta

  • A 调 acquireOrJoin(...)

  • 状态机创建一条新 meta,状态是 PENDING

  • A 成为 owner,拿到 ownerToken=101

第二步:A 进入运行态

  • A 调 markRunning(...)

  • 状态从 PENDING 变成 RUNNING

  • heartbeat 开始定时续租

第三步:B 到达

  • B 再调 acquireOrJoin(...)

  • Redis 发现当前状态活跃且 heartbeat 新鲜

  • B 收到 FOLLOWER_WAIT

  • B 不调模型,只等待 terminal event

第四步:A 执行成功

  • A 先写 result

  • 再把状态推进到 SUCCEEDED

  • 通过 Redis Stream 通知 follower

第五步:B 回放结果

  • B 被唤醒后读到 status=SUCCEEDED

  • B 读取 result,反序列化后直接返回

如果 A 在运行中失败,则第四步会变成:

  • A 把状态写成 FAILED

  • 如果失败类型可重试,后续请求可以重新把它接管回 PENDING

  • 如果不可重试,后续请求会直接读到失败结论

这就是当前状态机最核心的运行方式。