Java 项目用 Single-flight 一招终结 AI 接口重复调用

📅 2026/7/27 3:11:37
Java 项目用 Single-flight 一招终结 AI 接口重复调用
最近在做一个 AI 评分项目突然发现一个诡异的现象同一份简历、同一个 prompt后端日志显示模型被调了七八次。查了半天才发现——前端做了乐观更新用户点一次保存触发了三次请求加上页面重渲染和重试机制最夸张的一次同一个打分请求同时打到了后端 12 次。12 次每次都走了一遍完整的模型调用链路真金白银。后来我用 Single-flight 模式重构了这段逻辑同一个请求不管来多少次后端只调一次模型——这就是本文要讲的并发去重方案。一、背景为什么你需要关心重复调用AI 项目中的重复调用——这是当前最疼的场景。AI 接口有几个特点调用贵按 token 计费、延迟高动辄 210 秒、有并发限流。同一个评分请求、追问生成请求、简历抽题请求短时间内被多个线程同时打到同一个实例时每次调用都在烧钱。更要命的是很多 AI 平台有 RPM每分钟请求数限制并发重复请求占掉了宝贵的配额真正需要调用的请求反而被限流了。传统项目中的并发穿透——缓存失效瞬间大量请求同时打到数据库缓存击穿同一个外部 API支付查询、物流追踪被多个线程重复调用用户连点三次导出按钮触发三次同样的重计算。这些场景的共同点是N 个线程在极短时间内对同一个 key 发起相同的调用但真正需要执行的只有一次。二、Single-flight 是什么Single-flight单飞模式一种并发去重设计模式同一时刻对同一个 key 的多个并发请求只让第一个请求起飞执行其余请求阻塞等待并共享同一个结果。你可以理解为「机场只有一个跑道同一架航班只飞一次所有乘客共享这趟航班」。这个名字最早来自 Go 语言的golang.org/x/sync/singleflight包但设计思想是语言无关的。在 Java 生态里没有官方实现——这正是本文要补上的。Single-flight 的原理非常直观核心就是四件事打标签给每个请求分配一个唯一 key比如resume_123或请求参数的哈希相同 key 的请求视为同一件事。查状态内部维护一个 mapkey 是请求标识value 是当前有没有人正在执行。新请求进来先查这个 map。分角色如果 map 里没有这个 key说明你是第一个——你来执行并把状态写进 map如果 map 里已经有了说明有人已经在跑了——你等着就行。发结果执行完成后把结果返回给所有等待者然后把 key 从 map 里清掉。下次再有同样的请求重复上面的流程。用一个场景来感受一下你和几个同事同时走进一家餐厅都想点同一份麻婆豆腐。第一个人喊了一份麻婆豆腐服务员在菜单上记了一笔写进 map后面几个人听到之后就不喊了安心等着。后厨做了一份服务员端上来给所有人分了——而不是每个人喊一次、后厨做五份。先澄清一个容易混淆的点Single-flight ≠ 缓存。缓存是上次点过了这次直接加热端上来Single-flight 是别同时喊五次后厨做一次就够了。两者经常配合——缓存在前挡正常读Single-flight 在后挡并发穿透。三、单体版在 Java 中实现 Single-flight核心数据结构只有一个ConcurrentHashMapKey, CompletableFutureValue。key 是请求的去重标识比如resume_123value 是一个 Future第一个到达的线程负责创建 Future 并执行实际逻辑后续线程发现 Future 已存在就直接get()等待。下面看完整的单体版实现import java.util.concurrent.*; public class SingleFlightK, V { // 本地内存里的飞行中请求表。 // key请求的唯一标识比如 resume_123 // value一个正在执行或已执行完的 Future // 等待者通过它拿到最终结果。 private final ConcurrentHashMapK, CompletableFutureV inflight new ConcurrentHashMap(); /** * 对同一个 key 的并发请求只执行一次 loader其余等待并共享结果。 * * param key 去重标识例如对简历评分时传 resume_123 * param loader 真正要执行的逻辑比如调 AI 模型、查数据库 * param V 返回值类型支持任意类型 * return loader 的执行结果 */ public V execute(K key, CallableV loader) throws Exception { // 1. 快速路径先看看是不是已经有人在执行同一个 key 了。 // 如果已经有人在跑直接等它的结果不走后面的注册流程。 CompletableFutureV existing inflight.get(key); // get 不会加锁只是一个 volatile 读非常轻量。 if (existing ! null) { return existing.get(); // future.get() 会阻塞当前线程直到 owner 执行完。 // 如果 owner 已经执行完了get() 会立刻返回缓存的结果。 } // 2. 走到这里说明当前没有任何线程在执行这个 key。 // 创建一个新的 CompletableFuture 作为占位符 // 并通过 putIfAbsent 原子性地注册到 inflight 中。 CompletableFutureV future new CompletableFuture(); // 刚创建的 future 处于未完成状态 // 后面哪个线程是 owner 谁负责调用 complete 把它点亮。 CompletableFutureV 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而不是computeIfAbsentcomputeIfAbsent的问题不在于锁——实际上 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 模型四、单体版的能力边界与设计缺陷单体版在单实例场景下工作得很好但它有明确的边界1只在同一个 JVM 内有效。如果你部署了 3 个实例每个实例各自有一个SingleFlight实例互不可见。同一个 key 的请求如果被负载均衡分到了不同实例每个实例都会独立执行一次——这正是单体版最大的局限。2没有超时保护。如果 loader 里调 AI 接口卡住了 30 秒所有等待的线程也会卡 30 秒。生产环境里你需要给future.get()加上超时参数。3失败会传播给所有等待者。如果 loader 抛异常completeExceptionally会让所有future.get()的线程都收到同一个异常。在某些场景下这可能不是你想要的行为——你可能希望某一个等待者重试。4key 的粒度需要仔细设计。如果 key 太粗比如所有评分请求用同一个 key不同简历的请求会被错误合并如果 key 太细比如带上时间戳去重就失效了。5没有结果缓存。Single-flight 只管并发去重不管结果复用。如果 AI 的评分结果在 5 分钟内不会变你应该在外面套一层缓存比如 Caffeine而不是反复调 Single-flight。五、分布式版跨实例的请求合并单体版在单 JVM 内够用一上多实例就露馅。分布式 Single-flight 要解决的问题是不管请求落到哪个实例同一个 key 全局只执行一次。思路很直接用一个所有实例都能访问的中心化存储来做协调——在 Java 技术栈里Redis 是最自然的选择。核心流程分三步实例收到请求后先尝试在 Redis 里 SETNX 一个锁 keysf:lock:{key}拿到锁的实例负责执行 loader执行完后把结果写入 Redissf:result:{key}并通过 Pub/Sub 通知其他等待者没拿到锁的实例订阅 Redis Pub/Sub 频道等待结果通知超时则兜底轮询这里涉及两个关键的 Redis 原语SETNXSET if Not eXistsRedis 的原子命令仅当 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 ConcurrentHashMapString, CompletableFutureString 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, SupplierString loader) throws Exception { // 第一层本地去重 // 先查本地 inflight 表把同实例内的并发请求拦住 // 避免每个线程都跑去 Redis 抢锁浪费网络 IO 和 CPU。 CompletableFutureString local localCalls.get(key); // get 是 volatile 读不加锁非常轻量。 if (local ! null) { // 已有同实例线程在执行直接等它的结果。 return local.get( WAIT_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); // 带超时的 get防止 owner 卡死导致等待者永远挂住。 } // 没有本地在执行的记录尝试注册。 CompletableFutureString future new CompletableFuture(); // 刚创建的 future 处于未完成状态 // 谁注册成功谁负责执行完再 complete。 CompletableFutureString 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, SupplierString loader, CompletableFutureString future) throws Exception { // 拼接 Redis key用前缀区分不同用途避免冲突。 String lockKey sf:lock: key; // 分布式锁 key谁 SETNX 成功谁就是全局 owner。 String resultKey sf:result: key; // 结果 keyowner 执行完把结果写到这里 // 其他实例的等待者通过读这个 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, CompletableFutureString 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保证同实例内不再重复竞争。整体架构看这张图更直观