最近在做一个 AI 评分项目,突然发现一个诡异的现象:同一份简历、同一个 prompt,后端日志显示模型被调了七八次。查了半天才发现——前端做了乐观更新,用户点一次保存触发了三次请求,加上页面重渲染和重试机制,最夸张的一次,同一个打分请求同时打到了后端 12 次。

12 次,每次都走了一遍完整的模型调用链路,真金白银。后来我用 Single-flight 模式重构了这段逻辑,同一个请求不管来多少次,后端只调一次模型——这就是本文要讲的并发去重方案。

一、背景:为什么你需要关心重复调用

AI 项目中的重复调用——这是当前最疼的场景。AI 接口有几个特点:调用贵(按 token 计费)、延迟高(动辄 2~10 秒)、有并发限流。同一个评分请求、追问生成请求、简历抽题请求,短时间内被多个线程同时打到同一个实例时,每次调用都在烧钱。更要命的是,很多 AI 平台有 RPM(每分钟请求数)限制,并发重复请求占掉了宝贵的配额,真正需要调用的请求反而被限流了。

传统项目中的并发穿透——缓存失效瞬间大量请求同时打到数据库(缓存击穿),同一个外部 API(支付查询、物流追踪)被多个线程重复调用,用户连点三次"导出"按钮触发三次同样的重计算。

这些场景的共同点是:N 个线程在极短时间内对同一个 key 发起相同的调用,但真正需要执行的只有一次

二、Single-flight 是什么

Single-flight(单飞模式):一种并发去重设计模式,同一时刻对同一个 key 的多个并发请求,只让第一个请求"起飞"执行,其余请求阻塞等待并共享同一个结果。你可以理解为「机场只有一个跑道,同一架航班只飞一次,所有乘客共享这趟航班」。

这个名字最早来自 Go 语言的 golang.org/x/sync/singleflight 包,但设计思想是语言无关的。在 Java 生态里没有官方实现——这正是本文要补上的。

Single-flight 的原理非常直观,核心就是四件事:

  1. 打标签:给每个请求分配一个唯一 key(比如 "resume_123" 或请求参数的哈希),相同 key 的请求视为"同一件事"。
  2. 查状态:内部维护一个 map,key 是请求标识,value 是"当前有没有人正在执行"。新请求进来先查这个 map。
  3. 分角色:如果 map 里没有这个 key,说明你是第一个——你来执行,并把状态写进 map;如果 map 里已经有了,说明有人已经在跑了——你等着就行。
  4. 发结果:执行完成后,把结果返回给所有等待者,然后把 key 从 map 里清掉。下次再有同样的请求,重复上面的流程。

用一个场景来感受一下:你和几个同事同时走进一家餐厅,都想点同一份麻婆豆腐。第一个人喊了"一份麻婆豆腐!"服务员在菜单上记了一笔(写进 map),后面几个人听到之后就不喊了,安心等着。后厨做了一份,服务员端上来给所有人分了——而不是每个人喊一次、后厨做五份。

先澄清一个容易混淆的点:Single-flight ≠ 缓存。缓存是"上次点过了,这次直接加热端上来";Single-flight 是"别同时喊五次,后厨做一次就够了"。两者经常配合——缓存在前挡正常读,Single-flight 在后挡并发穿透。

三、单体版:在 Java 中实现 Single-flight

核心数据结构只有一个:ConcurrentHashMap<Key, CompletableFuture<Value>>。key 是请求的去重标识(比如 "resume_123"),value 是一个 Future,第一个到达的线程负责创建 Future 并执行实际逻辑,后续线程发现 Future 已存在就直接 get() 等待。

下面看完整的单体版实现:


import java.util.concurrent.*; public class SingleFlight<K, V> { // 本地内存里的"飞行中请求表"。 // key:请求的唯一标识,比如 "resume_123" // value:一个正在执行(或已执行完)的 Future, // 等待者通过它拿到最终结果。 private final ConcurrentHashMap<K, CompletableFuture<V>> inflight = new ConcurrentHashMap<>(); /** * 对同一个 key 的并发请求只执行一次 loader,其余等待并共享结果。 * * @param key 去重标识,例如对简历评分时传 "resume_123" * @param loader 真正要执行的逻辑,比如调 AI 模型、查数据库 * @param <V> 返回值类型,支持任意类型 * @return loader 的执行结果 */ public V execute(K key, Callable<V> loader) throws Exception { // 1. 快速路径:先看看是不是已经有人在执行同一个 key 了。 // 如果已经有人在跑,直接等它的结果,不走后面的注册流程。 CompletableFuture<V> existing = inflight.get(key); // get 不会加锁,只是一个 volatile 读,非常轻量。 if (existing != null) { return existing.get(); // future.get() 会阻塞当前线程,直到 owner 执行完。 // 如果 owner 已经执行完了,get() 会立刻返回缓存的结果。 } // 2. 走到这里说明当前没有任何线程在执行这个 key。 // 创建一个新的 CompletableFuture 作为"占位符", // 并通过 putIfAbsent 原子性地注册到 inflight 中。 CompletableFuture<V> future = new CompletableFuture<>(); // 刚创建的 future 处于未完成状态, // 后面哪个线程是 owner 谁负责调用 complete 把它点亮。 CompletableFuture<V> old = inflight.putIfAbsent(key, future); // putIfAbsent 是原子的"检查-设置"操作: // - 如果 key 不存在 → 写入并返回 null // - 如果 key 已存在 → 不写入,返回已存在的 value // 依靠这个原子操作来保证:并发场景下只有一个线程能注册成功。 if (old != null) { // 3. putIfAbsent 返回了非 null,说明在"get 看到 null" // 和"putIfAbsent 写入"之间,另一个线程抢先注册了。 // 那我就不再自己执行了,直接等那个线程的结果。 return old.get(); // old 就是抢先者创建的 future,等它就行。 } // 4. putIfAbsent 返回了 null,说明当前线程注册成功。 // 当前线程就是 owner,负责执行真正的 loader。 try { V result = loader.call(); // 真正执行调用方传进来的逻辑,比如调 AI 模型获取评分。 // 这个操作可能耗时几秒,但完全不在 ConcurrentHashMap 的锁内。 future.complete(result); // 把执行成功的结果写入 future。 // 此时所有在 future.get() 上阻塞的线程都会被唤醒, // 并拿到同一份 result。 return result; // owner 自己也返回这份结果。 } catch (Exception e) { // 如果执行过程中抛了异常,也要写进 future, // 否则等待者会一直挂住。 future.completeExceptionally(e); // completeExceptionally 会让所有 future.get() 的线程 // 抛出 ExecutionException,感知到同样的失败。 throw e; // owner 自己也要把异常继续往上抛。 } finally { // 不管成功还是失败,执行完一定要把 key 从 inflight 中移除。 // 否则:下次再有同一个 key 的请求进来,第一步 get 会发现 // 已有 future(虽然已经完成了),直接 get() 拿到旧结果—— // 就变成了一个不受控的缓存,违背了"执行完即清理"的语义。 inflight.remove(key); } } }

代码不长,但有几个设计细节值得展开:

为什么用 putIfAbsent 而不是 computeIfAbsent computeIfAbsent 的问题不在于锁——实际上 Java 8 之后 ConcurrentHashMap 内部已经没有分段锁了,mapping function 在锁外执行。真正的问题是:computeIfAbsent 不告诉你"是你插进去的,还是别人已经插过了",你没法判断自己是不是 owner。而 putIfAbsent 的返回值天然区分了这两种情况——返回 null 说明你注册成功你是 owner,返回非 null 说明别人抢先了你等着就行。

为什么在 finally 里 remove 不管 loader 成功还是抛异常,都要把 key 从 inflight 中移除,否则下一个请求过来发现 key 还在,会永远等在一个已经完成的 Future 上——虽然 get() 能正常返回,但这会导致 map 无限膨胀,内存泄漏。

看一张时序图把流程串起来——三个线程同时调用 execute("resume_123"),只有线程 1 真正调了 AI 模型:

Single-flight 时序图:三个线程并发调用同一个 key,只有线程1实际执行 AI 调用,其余线程共享结果

四、单体版的能力边界与设计缺陷

单体版在单实例场景下工作得很好,但它有明确的边界:

1)只在同一个 JVM 内有效。 如果你部署了 3 个实例,每个实例各自有一个 SingleFlight 实例,互不可见。同一个 key 的请求如果被负载均衡分到了不同实例,每个实例都会独立执行一次——这正是单体版最大的局限。

2)没有超时保护。 如果 loader 里调 AI 接口卡住了 30 秒,所有等待的线程也会卡 30 秒。生产环境里你需要给 future.get() 加上超时参数。

3)失败会传播给所有等待者。 如果 loader 抛异常,completeExceptionally 会让所有 future.get() 的线程都收到同一个异常。在某些场景下这可能不是你想要的行为——你可能希望某一个等待者重试。

4)key 的粒度需要仔细设计。 如果 key 太粗(比如所有评分请求用同一个 key),不同简历的请求会被错误合并;如果 key 太细(比如带上时间戳),去重就失效了。

5)没有结果缓存。 Single-flight 只管并发去重,不管结果复用。如果 AI 的评分结果在 5 分钟内不会变,你应该在外面套一层缓存(比如 Caffeine),而不是反复调 Single-flight。

五、分布式版:跨实例的请求合并

单体版在单 JVM 内够用,一上多实例就露馅。分布式 Single-flight 要解决的问题是:不管请求落到哪个实例,同一个 key 全局只执行一次

思路很直接:用一个所有实例都能访问的中心化存储来做协调——在 Java 技术栈里,Redis 是最自然的选择。

核心流程分三步:

  1. 实例收到请求后,先尝试在 Redis 里 SETNX 一个锁 key(sf:lock:{key}
  2. 拿到锁的实例负责执行 loader,执行完后把结果写入 Redis(sf:result:{key}),并通过 Pub/Sub 通知其他等待者
  3. 没拿到锁的实例,订阅 Redis Pub/Sub 频道等待结果通知,超时则兜底轮询

这里涉及两个关键的 Redis 原语:

SETNX(SET if Not eXists):Redis 的原子命令,仅当 key 不存在时才设置值,存在则不做任何操作。你可以理解为「第一个签到的人占住位置,后来的人看到已经有人签到了就自觉排队」。

Redis Pub/Sub(发布订阅):Redis 内置的消息广播机制,发布者向频道推送消息,所有订阅该频道的客户端实时收到。你可以理解为「广播喇叭——有结果了喊一声,所有等着的人都能听到」。

下面看实现:


import org.springframework.data.redis.core.StringRedisTemplate; import java.time.Duration; import java.util.concurrent.*; public class DistributedSingleFlight { // Redis 客户端,用来做分布式锁和结果传递。 private final StringRedisTemplate redis; // 第一层:本地"飞行中请求表"。 // 同实例内的并发先去这里去重,避免每个线程都去 Redis 抢锁。 // key:请求唯一标识,比如 "resume_123" // value:正在执行的 Future,等待者通过它拿结果。 private final ConcurrentHashMap<String, CompletableFuture<String>> localCalls = new ConcurrentHashMap<>(); // 分布式锁的 TTL,防止拿到锁的实例挂了导致锁永不释放。 // 30 秒足够覆盖绝大多数 AI 接口响应时间。 private static final Duration LOCK_TTL = Duration.ofSeconds(30); // 等待别人执行结果的超时时间。 // 设得比 LOCK_TTL 略短,避免等一个可能已经死掉的 owner。 private static final Duration WAIT_TIMEOUT = Duration.ofSeconds(25); public DistributedSingleFlight(StringRedisTemplate redis) { this.redis = redis; } /** * 对同一个 key 的并发请求,全局(跨实例)只执行一次 loader, * 其余请求等待并共享结果。 * * @param key 去重标识,比如 "resume_123" * @param loader 真正要执行的逻辑,比如调 AI 模型 * @return loader 的执行结果 */ public String execute(String key, Supplier<String> loader) throws Exception { // ===== 第一层:本地去重 ===== // 先查本地 inflight 表,把同实例内的并发请求拦住, // 避免每个线程都跑去 Redis 抢锁,浪费网络 IO 和 CPU。 CompletableFuture<String> local = localCalls.get(key); // get 是 volatile 读,不加锁,非常轻量。 if (local != null) { // 已有同实例线程在执行,直接等它的结果。 return local.get( WAIT_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); // 带超时的 get,防止 owner 卡死导致等待者永远挂住。 } // 没有本地在执行的记录,尝试注册。 CompletableFuture<String> future = new CompletableFuture<>(); // 刚创建的 future 处于未完成状态, // 谁注册成功谁负责执行完再 complete。 CompletableFuture<String> old = localCalls.putIfAbsent(key, future); // putIfAbsent 是原子操作:同实例内多个线程同时走到这里, // 只有第一个线程能写入成功(返回 null),其余拿到同一个 future。 if (old != null) { // 另一个同实例线程抢先注册了,等它的结果。 return old.get( WAIT_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); } // 当前线程是同实例内的 owner,负责去 Redis 抢全局执行权。 try { return doExecute(key, loader, future); } finally { // 不管成功还是失败,执行完后从本地表移除, // 否则同实例后续请求会一直复用一个已完成的 future。 localCalls.remove(key); } } /** * 第二层:Redis 分布式协调。 * 通过 SETNX 抢全局执行权,抢到的执行,没抢到的等结果。 */ private String doExecute(String key, Supplier<String> loader, CompletableFuture<String> future) throws Exception { // 拼接 Redis key,用前缀区分不同用途避免冲突。 String lockKey = "sf:lock:" + key; // 分布式锁 key:谁 SETNX 成功谁就是全局 owner。 String resultKey = "sf:result:" + key; // 结果 key:owner 执行完把结果写到这里, // 其他实例的等待者通过读这个 key 拿到结果。 String notifyChannel = "sf:notify:" + key; // 通知频道:owner 写完结果后发一条 Pub/Sub 消息, // 其他实例的等待者收到消息后立刻去读 resultKey。 // 尝试获取分布式执行权。 Boolean acquired = redis.opsForValue() .setIfAbsent(lockKey, "1", LOCK_TTL); // SETNX + EXPIRE 的原子组合: // - 如果 lockKey 不存在 → 写入 "1" 并设 30 秒 TTL,返回 true // - 如果 lockKey 已存在 → 什么都不做,返回 false // TTL 是兜底:万一 owner 挂了没删锁,30 秒后自动释放。 if (Boolean.TRUE.equals(acquired)) { // ===== 拿到全局执行权,我是 owner ===== try { String result = loader.get(); // 真正执行调用方的逻辑,比如调 AI 模型。 // 这一步可能耗时几秒,但完全在 Redis 锁的 TTL 范围内。 future.complete(result); // 先叫醒本地等待者(同实例内 other 线程)。 redis.opsForValue().set( resultKey, result, Duration.ofMinutes(5)); // 把结果写入 Redis,其他实例的等待者轮询时会读到。 redis.convertAndSend(notifyChannel, result); // 发一条 Pub/Sub 通知,其他实例的等待者收到后 // 立刻去读 resultKey,不用干等到下一次轮询。 return result; // owner 自己返回结果。 } catch (Exception e) { // 执行失败也要通知等待者,不能让他们干等。 future.completeExceptionally(e); // 本地等待者会收到 ExecutionException。 throw e; // owner 自己也要感知异常。 } finally { // 不管成功还是失败,一定要释放分布式锁。 // 否则其他实例的请求会一直认为有人在执行。 redis.delete(lockKey); } } else { // ===== 没拿到全局执行权,等别人执行完 ===== return waitForResult(key, resultKey, future); } } /** * 等待全局 owner 执行完,通过"先查一次 + Pub/Sub 通知 + 轮询兜底" * 三级策略获取结果。 */ private String waitForResult(String key, String resultKey, CompletableFuture<String> future) throws Exception { // 1. 先查一次:owner 可能刚执行完,结果已经在 Redis 里了。 String cached = redis.opsForValue().get(resultKey); if (cached != null) { future.complete(cached); // 叫醒本地其他等待者,它们还在等同一个 future。 return cached; } // 2. 订阅 Pub/Sub 通知(省略了具体的 MessageListener 注册代码)。 // 实际项目中可以用 Spring 的 RedisMessageListenerContainer, // 收到 notifyChannel 的消息后去读 resultKey 并 complete future。 // 这里简化为:由 Pub/Sub 监听器回调触发后续逻辑, // 同时下面第 3 步的轮询作为兜底。 // 3. 轮询兜底:Pub/Sub 不保证送达 long deadline = System.currentTimeMillis() + WAIT_TIMEOUT.toMillis(); // 计算截止时间,到点还没拿到结果就抛超时异常。 while (System.currentTimeMillis() < deadline) { cached = redis.opsForValue().get(resultKey); // 每 100ms 查一次,对 Redis 压力很小。 if (cached != null) { future.complete(cached); return cached; } Thread.sleep(100); // 100ms 间隔是轮询延迟和 Redis 压力的折中。 } // 等了 WAIT_TIMEOUT 还没拿到结果,说明 owner 可能挂了。 throw new TimeoutException( "等待 Single-flight 结果超时: " + key); // 调用方可以 catch 这个异常后决定重试还是降级。 } }

分布式版在单体版的基础上加了两层:Redis 分布式锁决定谁执行,本地 CompletableFuture 保证同实例内不再重复竞争。整体架构看这张图更直观:

Logo

欢迎加入 MCP 技术社区!与志同道合者携手前行,一同解锁 MCP 技术的无限可能!

更多推荐