← 返回笔记

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。这个类主要做两件事:

  1. 根据 stage、sessionId、题号、答案内容或文件 URL 生成稳定的 singleFlightKey

  2. 把真正的 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_NEWOWNER_TAKEOVER,就进入 owner 执行路径。

?????????????????????????? 完整流程如下:

  1. 调用 markRunning(...) 把状态从 PENDING 推进到 RUNNING

  2. 构造 FlightOwnerContext

  3. 调用 FlightHeartbeatManager.start(...) 启动周期 heartbeat

  4. 执行 supplier,也就是进入 AiCallGuardService.execute(...)

  5. AiCallGuardService 内部继续套上超时、Bulkhead、CircuitBreaker、Retry 等治理能力

  6. 最终通过 InterviewAiInvoker.doChat(...)XingChenAIClient.chat(...)

  7. 收到模型完整响应后,交给 FlightResultSerializer.serialize(...) 做压缩和 Base64 编码

  8. storeResult(...) 落 Redis 结果 hash

  9. finishSuccess(...) 把 meta 状态改成 SUCCEEDED

  10. FlightNotificationService.publish(...) 往 Redis Stream 发成功事件

  11. 把结果写入本地 FlightReplayLocalCache

  12. 停止 heartbeat,并把结果返回给业务层

这里真正触发第三方模型请求的,始终只有 owner 节点一个。这里看不懂不要紧,可以看一看这一片文档:

第七步:follower 节点等待并复用结果

如果当前节点拿到的是 FOLLOWER_WAIT,说明已经有别的节点在执行同一请求,这时它不会再去调模型,而是进入等待流程。

等待流程的逻辑是:

  1. 先再尝试一次 tryReadSuccessReplay(...)

  2. 读取 Redis meta,看是不是已经变成 SUCCEEDED

  3. 如果 meta 已经是 FAILEDretryable=false,直接抛错,不再等

  4. FlightNotificationService.waitForTerminalEvent(...),基于 Redis Stream 做阻塞等待

  5. 同时配合低频轮询兜底,防止只靠 Stream 导致遗漏

  6. 一旦发现成功结果,就读取 result hash

  7. 通过 FlightResultSerializer.deserialize(...) 反序列化并校验 checksum

  8. 放入本地 L1 缓存

  9. 返回给业务层

这就是“复用”真正发生的地方: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. 建议阅读顺序

如果要从代码角度顺着看,建议按下面顺序阅读:

  1. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/answer/InterviewEvaluationService.java

  2. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/answer/InterviewFollowUpService.java

  3. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/extraction/InterviewQuestionExtractionService.java

  4. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/flow/demeanor/InterviewDemeanorService.java

  5. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/shared/InterviewAiInvoker.java

  6. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/DistributedInterviewAiSingleFlightService.java

  7. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightCoordinatorRepository.java

  8. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightNotificationService.java

  9. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightHeartbeatManager.java

  10. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/FlightResultSerializer.java

  11. admin/src/main/java/com/hewei/hzyjy/xunzhi/interview/application/guard/AiCallGuardService.java

  12. 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 继续兜底。