5864 字
约 19 分钟
28
用 Single-flight 合并 AI 并发请求

从单机到分布式:用 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。

它会:

  1. 创建一个新的执行任务;
  2. 将任务注册到正在执行的请求表;
  3. 真正调用 AI;
  4. 保存最终结果或异常。

这个请求通常被称为:

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 执行链路建立一个统一、稳定、可治理的入口。

用 Single-flight 合并 AI 并发请求
http://www.clxhxhhr.top/posts/125/
作者
clxstart
发布于
2026-07-17
许可协议
CC BY-NC-SA 4.0
评论
0 条
还没有评论,先写一条吧。
文章目录
目录