从单机到分布式:用 Single-flight 合并 AI 并发请求
在 AI 面试、内容生成、智能评分等业务中,一次模型调用通常意味着明显的时间成本和费用成本。
当用户重复点击、前端自动重试,或者多个线程同时处理同一个任务时,同一个请求可能在几秒内被重复发送给模型。这样不仅会浪费 Token,还可能因为模型输出具有一定随机性,让同一份输入产生多份不同结果。
为了解决这个问题,我们在项目中引入了 Single-flight。
一句话总结
Single-flight 是一种并发请求合并技术:多个线程同时请求同一个资源时,只允许其中一个线程真正执行,其余线程等待并共享同一份结果。
构建的 JVM 内短生命周期请求复用器。
它可以解决同一应用实例内的 AI 请求重复执行问题,降低模型调用成本、线程池压力和结果抖动风险。
但它只能在单个 JVM 内生效,无法跨节点共享状态、进行故障接管或回放结果。因此,当系统演进为多实例部署后,还需要进一步升级为分布式 Single-flight。
一、什么是 Single-flight
假设线程 A 和线程 B 几乎同时发起了相同的 AI 评分请求,请求标识都是 k1。
没有 Single-flight 时:
线程 A ──→ 调用 AI
线程 B ──→ 调用 AI
系统会真实调用两次模型。
有 Single-flight 时:
线程 A ──→ 调用 AI ──→ 得到结果
↑
线程 B ──→ 等待 ──────┘
只有线程 A 真正调用 AI。
线程 B 不再执行重复操作,而是等待线程 A 完成,然后直接复用线程 A 的结果。
因此,Single-flight 可以被理解为:
正在执行中的请求去重器。
二、Single-flight 和缓存有什么不同
Single-flight 很容易被误解成缓存,但两者解决的问题并不相同。
缓存解决的是历史结果复用
例如:
第一次请求
→ 查询数据库
→ 将结果写入缓存
第二次请求
→ 直接读取缓存
即使两个请求间隔了十分钟,只要缓存没有过期,第二个请求仍然可以复用第一次的结果。
Single-flight 解决的是并发执行复用
例如:
请求 A 正在执行
请求 B 同时到达
→ B 等待 A
当 A 执行结束后,这次 Flight 通常就可以被移除。
所以二者的关注点不同:
| 机制 | 主要解决的问题 |
|---|---|
| 缓存 | 已经完成的结果能否继续复用 |
| Single-flight | 正在执行的相同请求是否需要重复执行 |
| 幂等控制 | 同一个业务操作是否会产生重复副作用 |
在实际系统中,这三种机制经常需要组合使用。
三、Single-flight 的核心原理
Single-flight 的运行过程可以分为四步。
1. 为请求生成唯一 Key
系统首先要判断哪些请求属于“同一个请求”。
例如,一个 AI 评分请求可以使用以下信息生成请求指纹:
面试记录 ID
+ 问题 ID
+ 用户答案
+ 评分模型
+ Prompt 版本
+ 业务阶段
然后对这些字段进行规范化和哈希:
SHA-256(normalized request)
只有 Key 相同的请求,才允许共享执行结果。
这是整个系统最重要的基础。
如果 Key 过于宽松,不同请求可能错误地共享结果;如果 Key 过于严格,本来相同的请求又无法合并。
2. 保存正在执行的请求
Single-flight 内部会维护一张“正在执行的请求表”:
Map<Key, FlightEntry>
其中:
Key
表示请求指纹
FlightEntry
表示正在执行的任务及其结果容器
3. 第一个请求成为执行者
当第一个请求到达时,系统中不存在对应的 Key。
它会:
- 创建一个新的执行任务;
- 将任务注册到正在执行的请求表;
- 真正调用 AI;
- 保存最终结果或异常。
这个请求通常被称为:
Owner
Leader
Flight Owner
4. 后续请求成为等待者
当其他线程携带相同 Key 到达时,会发现已有相同任务正在执行。
它们不会再次调用 AI,而是等待 Owner 完成。
请求 A ──→ 创建 Flight ──→ 调用 AI
请求 B ──→ 发现 Flight ──→ 等待
请求 C ──→ 发现 Flight ──→ 等待
如果 Owner 执行成功,所有等待者都会获得相同结果。
如果 Owner 执行失败,异常也会传递给所有等待者,避免请求一直处于等待状态。
四、单机 Single-flight 在项目中解决了什么
1. 减少同机重复调用 AI
项目中存在多个高成本 AI 场景,例如:
- 面试评分;
- 追问生成;
- 简历抽题;
- 神态分析;
- 内容总结;
- 报告生成。
只要多个请求最终落在同一个 JVM,并且请求 Key 相同,就可以被合并成一次真实的模型调用。
例如,同一时刻出现 10 个重复请求:
没有 Single-flight:
10 个请求 → 10 次 AI 调用
有 Single-flight:
10 个请求 → 1 次 AI 调用 + 9 个等待者
2. 降低瞬时资源压力
AI 调用通常还会占用:
- HTTP 连接;
- 业务线程;
- AI 调用线程池;
- 数据库连接;
- Token 配额;
- 下游服务并发额度。
Single-flight 可以从入口处收敛重复执行,避免相同任务重复消耗这些资源。
3. 降低 AI 结果抖动
大语言模型的输出并不一定完全确定。
同一份输入被并发调用两次,可能得到不同结果:
请求 A 返回:82 分
请求 B 返回:86 分
如果两个结果分别推进后续业务,可能引发:
- 数据互相覆盖;
- 状态冲突;
- 页面结果不一致;
- 重复发送消息;
- 业务流程重复推进。
Single-flight 让并发请求共享同一份模型输出,可以减少这种结果抖动。
4. 让等待者共享失败结果
如果 Owner 调用 AI 时失败,等待者不会一直挂起,也不会悄悄得到空结果。
Owner 的异常会被传播给所有等待者,让同一个 Flight 中的所有调用者观察到一致的失败结果。
五、TTL 和等待超时不是一回事
单机 Single-flight 中通常会涉及两个容易混淆的参数:
Flight TTL
等待超时
Flight TTL
TTL 表示一个 Flight 在本地内存中可以被认为有效多长时间。
超过 TTL 后,新请求可能不再复用旧 Flight,而是创建新的执行任务。
等待超时
等待超时表示一个等待者最多愿意等待 Owner 多长时间。
超过等待时间仍然没有结果,等待者就会结束等待并返回异常。
两者的含义不同:
TTL
控制 Flight 身份能够存在多久
等待超时
控制当前请求愿意等待多久
配置时必须考虑 AI 接口的真实耗时。
如果 AI 请求通常需要 20 秒,而 Flight TTL 只有 4 秒,那么旧任务还没有执行完成,新请求就可能创建新的 Flight,导致重复调用重新出现。
六、当前实现还带有短时间结果复用能力
严格意义上的 Single-flight 主要合并“正在执行中的请求”。
但如果一个已经完成的 Flight 在 TTL 内仍然保留,那么后续请求也可能直接读取刚刚完成的结果。
此时系统实际上结合了两种能力:
执行中的并发请求合并
+
极短时间内的结果回放
更准确地说,它属于:
带短 TTL 结果复用能力的本地 Single-flight。
这个设计对于以下场景非常有价值:
- 用户连续点击提交按钮;
- 前端快速自动重试;
- 网关短时间重新发送请求;
- 网络抖动造成重复提交。
但需要在设计文档中明确:系统不仅会合并正在执行的请求,也可能复用刚刚执行完成的结果。
七、单机版的能力边界
1. 只能在单 JVM 内生效
当前 Flight 状态保存在 JVM 本地内存中。
假设系统部署了两个实例:
Node A
Node B
两个节点会分别维护自己的 Flight 状态:
Node A:flightsA
Node B:flightsB
Node A 看不到 Node B 的执行状态,Node B 也看不到 Node A 的执行状态。
因此,同一个 Key 如果分别到达两个节点,两个节点仍然会各自执行一次 AI 调用。
2. 负载均衡会让相同请求落到不同节点
负载均衡是指在多个服务实例前增加一个统一入口,例如:
Nginx
API 网关
云负载均衡器
Kubernetes Service
负载均衡器会把请求分发给不同节点:
请求 1 → Node A
请求 2 → Node B
请求 3 → Node A
请求 4 → Node C
假设用户提交面试答案后,第一次请求落到了 Node A。
因为前端等待超时,又自动重试了一次,第二次请求落到了 Node B。
此时:
Node A:第一次看到 key=k1
Node B:也是第一次看到 key=k1
两个节点都会认为自己是第一个请求,并分别调用 AI。
所以:
单机 Single-flight 只能消除同一节点内的重复请求,无法消除跨节点重复请求。
3. 进程重启后状态全部丢失
Flight、执行状态和结果都保存在 JVM 堆内存中。
一旦发生:
- 服务重启;
- 容器重建;
- JVM 崩溃;
- 节点宕机;
- 滚动发布;
所有正在执行的 Flight 状态都会消失。
其他节点也无法知道原节点之前正在执行什么。
4. 无法进行跨节点故障接管
假设 Node A 获得了某个任务的执行权:
Node A:正在调用 AI
Node B:等待结果
如果 Node A 在执行过程中宕机,Node B 需要判断:
- Node A 是否真的已经死亡;
- 原任务是否仍然可能执行;
- 是否应该重新执行;
- 由谁成为新的 Owner;
- 如何防止两个 Owner 同时提交结果。
单机内存结构无法提供这些能力。
5. 无法跨节点共享结果
即使 Node A 已经成功完成 AI 请求,结果也只存在于 Node A 的内存中。
Node B 无法直接读取 Node A 中的 Future 或结果对象。
因此,分布式版本必须增加所有实例都可以访问的共享结果存储。
6. 缺少业务 Stage 级治理
不同 AI 业务的执行特征可能完全不同。
| Stage | 典型耗时 | 结果体积 | 重试特点 |
|---|---|---|---|
| 面试评分 | 5~15 秒 | 小 | 通常可以重试 |
| 追问生成 | 3~10 秒 | 中 | 通常可以重试 |
| 简历抽题 | 10~30 秒 | 较大 | 需要谨慎 |
| 神态分析 | 20~60 秒 | 较大 | 取决于输入 |
| 总报告生成 | 30~120 秒 | 很大 | 需要保证幂等 |
单机版本通常只有通用的:
- TTL;
- 等待超时;
- 清理阈值。
但它无法按照不同 Stage 配置:
- Owner 租约;
- 心跳频率;
- 结果保存时间;
- 结果压缩;
- 重试策略;
- 最大等待者数量;
- 是否允许故障接管;
- 是否启用本地 L1 结果复用。
八、单机版本仍需关注的细节
1. TTL 过期不代表原任务停止
假设 AI 调用需要 20 秒,而 TTL 只有 4 秒:
0 秒:Owner A 开始调用 AI
4 秒:Flight 过期
5 秒:新请求到达
5 秒:Owner B 创建新 Flight
此时 A 和 B 会同时调用 AI。
所以 TTL 更准确的含义是:
多久以后允许新请求放弃复用旧 Flight,并创建新的执行任务。
它并不是 Owner 的强制取消时间。
2. 等待者超时不代表 Owner 已停止
某个等待请求超时,只能说明:
当前等待者不愿意继续等待
并不代表:
Owner 已经停止执行
如果等待者超时后直接删除 Flight,而 Owner 仍在执行,那么后续请求可能创建新的 Flight,再次调用 AI。
更稳妥的策略包括:
- 等待者只退出,不删除 Owner 的 Flight;
- 仅由 Owner 负责清理 Flight;
- 后台统一清理异常状态;
- 结合租约和心跳判断 Owner 是否真正失活。
3. Key 必须包含业务类型和版本
如果不同业务错误地使用相同 Key,却期待不同类型的结果,就可能出现错误复用。
因此,Key 最好包含:
业务名称
+ Stage
+ Prompt 版本
+ 模型版本
+ 请求参数指纹
例如:
interview:score:v2:{requestHash}
interview:followup:v1:{requestHash}
resume:question:v3:{requestHash}
这样可以避免不同业务或不同版本之间发生错误共享。
4. 需要限制 Flight 数量
如果短时间出现大量不同 Key,本地 Flight 表可能持续增长。
因此需要考虑:
- 最大 Flight 数量;
- 超限拒绝策略;
- 按 Stage 限流;
- 定时清理任务;
- 活跃 Flight 数量监控;
- 内存使用告警。
九、为什么不能只增加一个分布式锁
升级为分布式版本时,最容易想到的方案是:
Redis SET NX
第一个节点获得锁后负责执行,其他节点没有获得锁就等待。
但只有分布式锁还不够。
分布式锁只能回答:
谁负责执行?
它无法完整回答:
- 任务当前是什么状态?
- Owner 是否仍然存活?
- 结果保存在哪里?
- 等待者如何获取结果?
- Owner 宕机后谁来接管?
- 旧 Owner 恢复后还能不能写入结果?
- 完成结果可以回放多久?
因此,一个完整的分布式 Single-flight 至少需要:
分布式执行权
+ 执行状态
+ 租约和心跳
+ 结果存储
+ 等待通知
+ 故障接管
十、分布式 Single-flight 的基本设计
一个较完整的分布式方案可以分为两层:
L1:JVM 本地 Single-flight
L2:分布式 Single-flight
请求首先经过本地层:
同一 JVM 内的重复请求
→ 使用本地内存快速合并
本地 Owner 再进入分布式层:
不同节点之间的重复请求
→ 使用 Redis 或其他共享组件合并
整体流程如下:
业务请求
↓
生成请求指纹
↓
JVM 本地 Single-flight
↓
尝试获得分布式执行权
├── 成功:成为 Owner,执行 AI
└── 失败:成为 Follower,等待共享结果
1. 保存分布式任务状态
可以在 Redis 或数据库中维护:
singleflight:{stage}:{key}
状态内容可以包括:
{
"status": "RUNNING",
"ownerId": "node-a",
"generation": 18,
"startedAt": 1710000000000,
"heartbeatAt": 1710000003000,
"expireAt": 1710000030000
}
任务状态可以设计为:
RUNNING
SUCCEEDED
FAILED
TIMEOUT
2. 使用租约而不是永久锁
Owner 获得的执行权不应该永久存在,而应该是一段有限时间的租约。
例如:
租约时间:30 秒
心跳间隔:5 秒
Owner 执行期间持续续租:
Node A 每 5 秒更新一次 heartbeat
如果心跳长时间没有更新,其他节点可以判断 Owner 可能已经失活,并尝试接管任务。
3. 将结果写入共享存储
Owner 执行成功后,需要把结果写入所有节点都能访问的共享存储:
singleflight:result:{stage}:{key}
结果可以包含:
{
"status": "SUCCEEDED",
"result": "...",
"completedAt": 1710000010000,
"expireAt": 1710000070000
}
其他节点不需要访问 Owner JVM 中的 Future,而是直接读取共享结果。
如果 AI 结果体积较大,可以考虑:
- GZIP 压缩;
- 保存到数据库;
- 保存到对象存储;
- Redis 只保存结果地址;
- 根据 Stage 设置不同的结果 TTL。
4. 等待者如何获得结果
Follower 可以采用三种方式等待结果。
定时轮询
每隔 200~500 毫秒查询一次任务状态
实现简单,但会增加 Redis 或数据库压力。
发布订阅
Owner 完成后发送通知:
singleflight-completed:{key}
Follower 收到通知后读取共享结果。
这种方式响应较快,但发布订阅消息可能丢失,因此不能只依赖通知。
通知与轮询结合
更稳妥的方案是:
通知用于快速唤醒
结果存储用于最终结果读取
轮询用于消息丢失兜底
5. 使用 Generation 防止旧 Owner 回写
假设 Node A 的租约已经过期,Node B 接管了任务。
但 Node A 并没有真正死亡,只是网络短暂抖动,之后它又完成了 AI 调用。
此时可能出现:
Node B 已成为新 Owner
Node A 仍然尝试写入旧结果
因此,每次任务被接管时都应该增加一个版本号:
generation = 18
generation = 19
Owner 提交结果时,必须校验自己持有的 Generation 是否仍然有效:
只有当前 Generation 的 Owner
才能提交最终结果
这种版本号也常被称为:
Fencing Token
它可以防止已经失去执行权的旧节点覆盖新节点的结果。
十一、按照业务 Stage 配置策略
分布式版本不应该只使用一组全局参数。
可以按照不同 AI 任务配置不同策略:
single-flight:
stages:
interview-score:
lease-timeout: 30s
heartbeat-interval: 5s
wait-timeout: 35s
result-ttl: 60s
resume-question:
lease-timeout: 90s
heartbeat-interval: 10s
wait-timeout: 100s
result-ttl: 10m
compress-result: true
final-report:
lease-timeout: 180s
heartbeat-interval: 15s
wait-timeout: 200s
result-ttl: 30m
不同 Stage 可以拥有不同的:
- 租约时间;
- 心跳间隔;
- 等待超时;
- 结果 TTL;
- 结果压缩;
- 重试策略;
- 最大并发数;
- 失败结果是否短暂保存;
- 是否允许自动接管。
十二、需要监控哪些指标
Single-flight 是否真正有效,不能只看代码逻辑,还需要监控实际运行效果。
建议至少记录:
singleflight_request_total
singleflight_owner_total
singleflight_follower_total
singleflight_hit_total
singleflight_miss_total
singleflight_wait_timeout_total
singleflight_owner_failure_total
singleflight_takeover_total
singleflight_active_flights
singleflight_wait_duration
singleflight_execution_duration
其中一个重要指标是合并率:
合并率 = Follower 请求数 / 总请求数
还可以统计节省的 AI 调用次数:
节省调用数 = 总业务请求数 - Owner 执行数
例如:
总业务请求数:1000
真实 AI 调用数:620
说明 Single-flight 合并了约 380 次重复调用。
这些数据还可以进一步换算为:
- Token 成本节省;
- 模型调用费用节省;
- 下游并发压力下降;
- 平均等待时间变化;
- 结果一致性提升。
十三、建议的演进路线
项目中的 Single-flight 可以分三个阶段演进。
阶段一:JVM 本地版
当前方案:
ConcurrentHashMap
+ CompletableFuture
+ TTL
+ 等待超时
目标:
- 快速解决单节点并发重复;
- 不引入额外中间件;
- 保持实现简单;
- 验证 Single-flight 的业务价值。
阶段二:本地加分布式协调
系统结构演进为:
本地 ConcurrentHashMap
↓
Redis 分布式租约
↓
共享结果存储
目标:
- 合并跨节点重复请求;
- 支持结果共享;
- 支持基础故障判断;
- 保留 JVM 内快速路径。
阶段三:业务级 AI 任务协调系统
进一步增加:
- Stage 级策略;
- Owner 心跳;
- Generation 和 Fencing Token;
- 故障接管;
- 结果压缩;
- 长任务状态查询;
- 异步通知;
- 失败重试;
- 审计记录;
- 成本统计。
此时,它已经不只是一个并发工具类,而是一个面向 AI 任务的分布式执行协调组件。
总结
Single-flight 的核心思想并不复杂:
多个并发请求执行同一件事时,只让一个请求真正执行,其余请求共享结果。
当前基于 ConcurrentHashMap + CompletableFuture 的实现,用较为克制的方式解决了单 JVM 内 AI 请求重复执行的问题。
它带来的价值包括:
- 减少重复模型调用;
- 降低 Token 和接口成本;
- 缓解线程池与下游压力;
- 减少同一输入产生多个不同结果;
- 让并发调用获得一致的成功或失败结果。
但它的能力边界也非常明确:
- 无法跨 JVM;
- 无法跨节点共享状态;
- 无法进行故障接管;
- 无法跨节点回放结果;
- 缺少租约、心跳和 Stage 级治理。
因此,当前单机版不是最终形态,而是一个合理的第一阶段:
单 JVM 请求合并
↓
跨节点执行协调
↓
分布式 AI 任务治理
先用简单可靠的本地实现验证业务价值,再逐步引入分布式租约、结果存储、心跳和故障接管,比一开始就构建复杂状态机更加稳妥。
对于高成本、长耗时、结果具有随机性的 AI 调用来说,Single-flight 不只是一次并发优化,更是在为整个 AI 执行链路建立一个统一、稳定、可治理的入口。