learning note
设计分布式 Single-flight 落地
完整梳理分布式 Single-flight 的落地设计,包括请求唯一 key、Redis 状态、owner、follower 和结果复用。
背景:
AI接口是比较宝贵的资源。因此我们希望实现减少并发重复调用AI的情况。比如同一个评分请求、追问生成请求、简历抽题请求,短时间内有多个线程同时打到同一个实例时,不应该每个线程都去调一次模型。
那针对这个问题,其实业内有完整的解决方案:Single-Flight(Single-flight是什么?(必看)。而我们本次设计的并不是单体Single-Flight,而是分布式Single-Flight。单体Single与分布式Single的区别(必看)
核心参与类以及职责:
| 类 | 作用 |
|---|---|
InterviewAiInvoker |
统一 AI 调用入口,负责构造 single-flight key,并把实际调用包装成统一 supplier,兄弟们把这个supplier看作是要执行的AI逻辑就行。 |
DistributedInterviewAiSingleFlightService |
集群级 single-flight 核心,实现 owner/follower 协调、回放、接管、降级 |
FlightCoordinatorRepository |
Redis 协调层,通过 Lua 脚本维护状态机、结果和心跳 |
FlightNotificationService |
用 Redis Stream 做 owner -> follower 的终态通知 |
FlightHeartbeatManager |
owner 的续租调度器,定期刷新运行态 heartbeat |
FlightResultSerializer |
对 AI 结果做压缩、Base64 编码、checksum 校验 |
FlightReplayLocalCache |
节点内 L1 回放缓存,减少重复读 Redis |
AiCallGuardService |
真正执行模型调用前的保护层,负责超时、隔离、熔断、重试 |
InterviewAiSingleFlightService |
JVM 内单机复用器,作为禁用分布式或 hybrid 降级时的本地兜底 |
InterviewAiSessionLockService |
对 extraction、demeanor 这类重任务做会话级互斥 |
整体架构思路:
当前实现不是单独增加一个缓存组件,而是把 AI 调用拆成四层:
1. 业务入口层:评分、追问、抽题、神态分析。
2. 重任务互斥层:InterviewAiSessionLockService对 extraction、demeanor 做 session 级重锁。
3. 集群级协调层:DistributedInterviewAiSingleFlightService + FlightCoordinatorRepository 负责 owner/follower 协调。
4. 调用保护层:AiCallGuardService 对真正的 owner 调用做超时、熔断、隔离、重试保护。
从架构上看,这种拆法的关键点在于:
-
重锁解决的是“同一会话的重任务不要并发”。
-
Single-flight解决的是“同一份 AI 请求不要重复调用”。
-
AI guard解决的是“真正要调模型时,如何安全地调”。
哪些业务入口会走这条链路
请先看这篇文章了解什么是Stage:什么是Stage
面试评分
-
入口类:
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/answer/InterviewEvaluationService.java -
关键方法:
evaluateAnswerByScorerAgent(...) -
调用方式:通过
InterviewAiInvoker.callAiSyncWithParameters(...)调 scorer workflow -
key 维度:
stage`` + sessionId + questionNumber + answerHash
这个场景会把“候选人回答、题目内容、简历上下文”组装成工作流参数,然后进入统一 AI 调用链。因为评分请求天然具备“同题同回答”的复用价值,所以 key 用题号和答案摘要来归一化。
追问问题生成
-
入口类:
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/answer/InterviewFollowUpService.java -
关键方法:
generateFollowUpQuestion(...)/invokeFollowUpWorkflow(...) -
调用方式:通过
InterviewAiInvoker.callAiSyncWithParameters(...)调 follow-up workflow -
key 维度:
stage + sessionId + currentQuestion + answerHash
这个场景的复用目标是:同一会话、同一道当前题、同一份回答,不要重复生成追问。
简历题目抽取
-
入口类:
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/extraction/InterviewQuestionExtractionService.java -
关键方法:抽取主流程里先
acquire(...)重锁,再callAiSyncWithFile(...) -
调用方式:上传简历文件后,按文件 URL 作为业务 key 进入统一 AI 调用链
-
key 维度:
stage + sessionId + fileUrl
这个场景比普通问答更重,因为前面还有文件上传、后面还有结构化落库,所以它不仅走分布式 Single-flight,还先走 session 级重锁。
神态分析
-
入口类:
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/demeanor/InterviewDemeanorService.java -
关键方法:
evaluateDemeanor(...) -
调用方式:图片上传后,通过
InterviewAiInvoker.callAiSyncWithFile(...)调用神态分析 workflow -
key 维度:
stage + sessionId + imageUrl
这个场景同样属于重任务,也先走 session 级重锁,再进入分布式 Single-flight。
工作流程
我们的后端会先把所有的AI请求归一化为一个稳定的key。比如AI评分请求带上stage + sessionId + questionNumber + answerHash,这样无论这个请求被负载均衡打到哪台机器,只要业务语义相同,所有节点算出来的 key 都一样。
请求进入DistributedInterviewAiSingleFlightService后,会先查一次本地的 L1 回放缓存;如果没命中,就去 Redis 里通过 Lua 脚本做一次原子抢占。抢占结果只会有几种:
-
要么当前节点成为 owner
-
要么当前节点成为 follower
-
要么直接发现历史成功结果可以复用
Owner的情况:
如果当前节点是 owner,它会先把状态推进到 RUNNING,然后启动 heartbeat 定时续租,表示我还活着、还在执行,不要被别人接管。接着它才真正去执行 supplier,也就是进入AiCallGuardService,最终调用模型。
调用成功后,会把结果做序列化、压缩、checksum 校验后写入 Redis,再把状态改成 SUCCEEDED,并通过 Redis Stream 通知正在等待的 follower。调用失败时,则会把状态改成 FAILED,同时标记失败类型和是否允许重试
如果执行过程中 owner 挂了、卡住了、长时间不发 heartbeat,别的节点会检测到它失活,然后基于新的 ownerToken 安全接管,避免旧 owner 恢复后把旧结果写回来,这就是 fencing token 的作用。
follower的情况:
如果当前节点是 follower,它不会真正调 AI,而是等待 owner 的结果。等待过程不是傻轮询,而是Redis Stream 阻塞等待 + 低频轮询兜底。
一旦发现 owner 成功了,就直接从 Redis 读取结果、反序列化后返回;如果发现是不可重试失败,就直接结束,不会继续重复打模型。
降级情况:
如果分布式等待或协调过程本身出现异常。并且当前模式允许降级,那么也会回退到 InterviewAiSingleFlightService。以单机内 single-flight 的方式继续兜底,保证系统在分布式协调链路不稳定时仍然具备基本的去重能力。
说白了就是大部分情况下走分布式Single-Flight,但是如果分布式服务出现bug,比如Redis挂了导致没办法选Owner了,那我们就用单机Single-Flight。
链路图
??????????????????????????
Owner/Follower 主流程时序图
Redis Meta 生命周期图
详细工作流程解释:
第一步:业务层构造参数和 key
所有业务场景最终都会收敛到 InterviewAiInvoker。这个类主要做两件事:
-
根据 stage、sessionId、题号、答案内容或文件 URL 生成稳定的
singleFlightKey -
把真正的 AI 调用封装成
supplier,交给分布式 Single-flight 服务处理
当前 key 生成规则有两类:
-
问答类::
stage+sessionId+questionNumber+answerHash -
文件类::
stage+sessionId+businessKey
这样设计的结果是:同一个业务语义的请求,即使来自不同节点,也会落到同一个 key 上,从而进入同一个 flight。
key具体是怎么构造的,请看这一篇文章:归一化Key的具体设计
第二步:重任务先走 session 级重锁
不是所有 stage 都先加Session重锁,当前只有两类重任务会这样做:
-
interview-extraction(简历题目抽取) -
interview-demeanor(神态分析)
原因很简单:这两个场景除了 AI 本身,还会伴随文件上传、图片上传、长耗时处理、后续结构化解析等逻辑。如果同一个会话并发触发多次,即使分布式 Single-flight 最终只保留一个 owner,也可能在前置重操作上浪费资源。所以这里先用 InterviewAiSessionLockService 做 session 级互斥,再进入 AI 调用链。
评分和追问则不走这层重锁,因为它们更轻量,也更适合直接依赖 single-flight 的请求级复用。
Session重锁是怎么构造的,可以看这篇文章:Session重锁的具体设计
第三步:进入分布式 Single-flight 主入口
InterviewAiInvoker 并不是直接调用模型,而是把 supplier 包装成下面这条链:
DistributedInterviewAiSingleFlightService.execute(stage, key, () -> aiCallGuardService.execute(stage, key, callable))
也就是说我们先执行DistributedInterviewAiSingleFlightService,再执行aiCallGuardService
两者关系
-
DistributedInterviewAiSingleFlightService 管“同一个集群的同一语义请求只被调用一次”
-
AiCallGuardService 管“如果真的要调用AI,这次调用怎么发得更稳、更安全”
所以这条链的意思就是:先由分布式 Single-Flight 判断谁有资格执行;只有那个真正有资格的 owner,才会进入 AiCallGuardService,然后再去调模型。
第四步:先查本地 L1 回放缓存
DistributedInterviewAiSingleFlightService 收到请求后,第一件事不是立刻访问 Redis,而是先查 FlightReplayLocalCache。
如果当前节点在短时间内已经回放过这个 key 的成功结果,就直接返回,不再打 Redis。这一层不是用来解决并发抢占的,而是用来降低 Redis 读压和缩短后续同 key 请求的响应路径。
注意,这个 L1 缓存不是所有 stage 都开:
-
interview-evaluation:开 -
interview-followup:开 -
interview-extraction:开 -
interview-demeanor:关
神态分析之所以默认不开,是因为它更偏媒体分析,结果复用窗口通常更谨慎。毕竟神态这个东西,每一秒的状态都不一样,很难存在复用的情况。
而LocalCache内部其实就是一个带自动过期时间的Map而已,没啥牛逼的。
??????????????????????????
第五步:Redis Lua 抢占 owner / follower 身份
如果 L1 没命中,服务会调用 FlightCoordinatorRepository.acquireOrJoin(...)。这一层通过 Lua 脚本在 Redis 上维护 flight 元数据。
?????????????????????????? 当前脚本逻辑可以概括成下面几种分支(涉及到对Redis状态机的修改了,所以先看这篇文章:Redis状态机设计思路):
状态机看完之后,再看这篇文章:Action分类,一定要看,不然你不理解下面是在干什么
分支 A:第一次请求
如果 Redis 里还没有 meta,就创建一条 PENDING 记录,并分配新的 ownerToken,返回 OWNER_NEW。
分支 B:已经成功过
如果状态已经是 SUCCEEDED,直接返回 REPLAY_SUCCESS,调用方后续去读结果即可。
分支 C:之前失败过
如果状态是 FAILED:
-
retryable = true:允许新节点接管,返回OWNER_TAKEOVER -
retryable = false:直接返回REPLAY_FAILURE,表示这个失败结果不应该被继续重试
分支 D:当前 owner 还活着
如果状态还在运行窗口内,并且 heartbeatAt 距当前时间没有超过 takeoverDetectMillis,说明 owner 还活着,这时当前节点成为 FOLLOWER_WAIT。
分支 E:owner 疑似失活
如果 heartbeat 已经长时间不更新,说明 owner 可能挂了、卡住了或者超时了,这时会重新分配 ownerToken,返回 OWNER_TAKEOVER。
第六步:owner 节点真正执行 AI 调用
如果当前节点拿到的是 OWNER_NEW 或 OWNER_TAKEOVER,就进入 owner 执行路径。
?????????????????????????? 完整流程如下:
-
调用
markRunning(...)把状态从PENDING推进到RUNNING -
构造
FlightOwnerContext -
调用
FlightHeartbeatManager.start(...)启动周期 heartbeat -
执行 supplier,也就是进入
AiCallGuardService.execute(...) -
AiCallGuardService内部继续套上超时、Bulkhead、CircuitBreaker、Retry 等治理能力 -
最终通过
InterviewAiInvoker.doChat(...)调XingChenAIClient.chat(...) -
收到模型完整响应后,交给
FlightResultSerializer.serialize(...)做压缩和 Base64 编码 -
调
storeResult(...)落 Redis 结果 hash -
调
finishSuccess(...)把 meta 状态改成SUCCEEDED -
调
FlightNotificationService.publish(...)往 Redis Stream 发成功事件 -
把结果写入本地
FlightReplayLocalCache -
停止 heartbeat,并把结果返回给业务层
这里真正触发第三方模型请求的,始终只有 owner 节点一个。这里看不懂不要紧,可以看一看这一片文档:
第七步:follower 节点等待并复用结果
如果当前节点拿到的是 FOLLOWER_WAIT,说明已经有别的节点在执行同一请求,这时它不会再去调模型,而是进入等待流程。
等待流程的逻辑是:
-
先再尝试一次
tryReadSuccessReplay(...) -
读取 Redis meta,看是不是已经变成
SUCCEEDED -
如果 meta 已经是
FAILED且retryable=false,直接抛错,不再等 -
调
FlightNotificationService.waitForTerminalEvent(...),基于 Redis Stream 做阻塞等待 -
同时配合低频轮询兜底,防止只靠 Stream 导致遗漏
-
一旦发现成功结果,就读取 result hash
-
通过
FlightResultSerializer.deserialize(...)反序列化并校验 checksum -
放入本地 L1 缓存
-
返回给业务层
这就是“复用”真正发生的地方:follower 并不是复用 Redis 状态,而是复用 owner 已经算出来的最终 AI 结果。
第八步:失败、接管和降级
如果 owner 在执行过程中出错,DistributedInterviewAiSingleFlightService 会先通过 classifyFailure(...) 把异常分类为:
-
TIMEOUT -
OVERLOAD -
PROVIDER -
VALIDATION -
UNEXPECTED
然后调用 finishFailure(...) 把 meta 改成 FAILED,并记录:
-
errorType -
errorCode -
retryable
后续就会分成两类:
- 可重试失败:后面的节点再次进来时,可以 `OWNER_TAKEOVER`
- 不可重试失败:后面的 follower 会直接收到失败,不再继续调模型
除此之外,当前配置是 mode: hybrid。这意味着如果分布式协调链路本身发生运行时异常,比如 Redis 协调出现问题,那么系统不会直接中断,而是会回退到 InterviewAiSingleFlightService,至少在单机内继续维持一次复用,作为可用性兜底。
11. 建议阅读顺序
如果要从代码角度顺着看,建议按下面顺序阅读:
-
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/answer/InterviewEvaluationService.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/answer/InterviewFollowUpService.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/extraction/InterviewQuestionExtractionService.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/demeanor/InterviewDemeanorService.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/shared/InterviewAiInvoker.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/DistributedInterviewAiSingleFlightService.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightCoordinatorRepository.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightNotificationService.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightHeartbeatManager.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightResultSerializer.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/AiCallGuardService.java -
admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/InterviewAiSingleFlightService.java
按这个顺序看,最容易把“业务入口 -> 统一调用 -> 集群协同 -> 模型执行 -> 结果回放”整条线串起来。
一句话总结:
业务服务先按 stage 和业务语义生成稳定 key,重任务先做 session 级互斥,然后进入分布式 Single-flight;集群里只有一个 owner 会真正经过 AI guard 去调模型,结果会被序列化后写入 Redis,并通过 Stream 通知 follower,其他节点只等待并复用结果;如果分布式协调异常,则在 hybrid 模式下回退到本地 single-flight 继续兜底。