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:
-
这个 key 第一次出现,
acquireOrJoin(...)抢占成功 -
之前是
FAILED且允许重试,新的 owner 接管成功 -
之前 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
也就是说,真正高频发生的是这几件事:
-
抢占成功,生成
PENDING -
owner 正式开跑,进入
RUNNING -
执行完成后进入
SUCCEEDED或FAILED -
如果失败可重试,再次接管,又重新回到
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 -
如果不可重试,后续请求会直接读到失败结论
这就是当前状态机最核心的运行方式。