一、引言那个让数据库瑟瑟发抖的瞬间想象这样一个场景凌晨三点一个热点缓存 key 恰好过期。下一秒成千上万的请求如潮水般涌向你的服务——每一个请求都发现缓存为空于是不约而同地冲向数据库。数据库连接池瞬间被占满CPU 飙升响应时间急剧恶化最终服务不可用。这就是缓存击穿Cache Stampede / Thundering Herd。互斥锁能解决吗可以但代价是所有请求串行化——1000 个请求就得排队等 1000 次性能损失惨重。而 SingleFlight正是为解决这个问题而生的请求合并利器。SingleFlight 的核心思想对于同一个 key 的多个并发请求只让第一个真正执行其余的阻塞等待待第一个完成后所有请求共享同一个结果。它的名字很形象——单次飞行一群鸟里只让一只飞出去其他的等着分享它的见闻。二、SingleFlight 是什么2.1 一句话定义SingleFlight 是一种重复函数调用抑制机制它将相同 key 的并发请求合并为一次实际调用所有请求共享该调用的结果。2.2 一个直观的比喻如果多个人同时想点外卖与其每个人都打开外卖 APP 下单不如大家凑在一起由一个人下单然后大家共享这份外卖。这就是 SingleFlight 的工作方式。2.3 适用场景缓存击穿防护热点缓存失效时防止大量请求穿透到数据库接口限流降级高并发下减少对下游服务的重复调用分布式锁优化配合分布式锁减少锁竞争任何昂贵操作的去重如 RPC 调用、慢 SQL 查询、AI 模型推理等三、SingleFlight 工作原理3.1 核心数据结构SingleFlight 的核心由两个结构体组成// Group管理所有请求的命名空间 type Group struct { mu sync.Mutex // 互斥锁保护 map 的并发安全 m map[string]*call // key → 正在执行的请求 } // call代表一个进行中或已完成的请求 type call struct { wg sync.WaitGroup // 用于阻塞等待结果 val interface{} // 执行结果 err error // 执行错误 dups int // 重复调用计数 chans []chan- Result // 等待结果的通道列表 }3.2为什么它能防缓存击穿假设 1000 个请求同时查询同一个已过期的缓存 key方案数据库查询次数响应方式无保护1000 次全部穿透数据库崩溃互斥锁1000 次串行排队等待RT 飙升SingleFlight1 次并发等待共享结果关键在于并发量不变但下游压力降为 1。四、Java 实现 SingleFlight4.1 基于 ConcurrentHashMap CompletableFuture 的实现这是最经典的 Java 实现方式使用ConcurrentHashMap管理进行中的请求CompletableFuture实现异步等待/** * SingleFlight - 请求合并器 * * 核心思想当多个并发请求使用相同的 key 时只让第一个请求真正执行 * 后续请求直接复用第一个请求的结果从而避免重复的昂贵操作如 DB 查询、RPC 调用。 * * 类比10 个人同时想查同一本书只让 1 个人去仓库取剩下 9 个人在原地等 * 等书取回来大家一起看。 * * param T 返回结果的类型 */ public class SingleFlightT { // 核心存储结构 /** * 存储正在进行的请求key → CompletableFutureT * * 为什么用 ConcurrentHashMap * 因为会有多个线程同时尝试 put key需要线程安全的 Map。 * * 为什么 key 对应的是 CompletableFuture而不是直接存结果 * 因为后续请求需要能够“等待”结果而不仅仅是拿到最终值。 * CompletableFuture 天然支持多个线程可以同时调用 future.get() * 一旦 future 完成成功或失败所有等待线程都会同时被唤醒。 */ private final ConcurrentHashMapString, CompletableFutureT ongoingRequests new ConcurrentHashMap(); // 自定义线程池可选默认使用 ForkJoinPool.commonPool() private final Executor executor; // 构造器 /** * 默认构造器使用公共 ForkJoinPool 执行 loader */ public SingleFlight() { this(ForkJoinPool.commonPool()); } /** * 自定义线程池构造器便于隔离资源和监控 */ public SingleFlight(Executor executor) { this.executor executor; } // 核心方法同步阻塞版 /** * 执行或等待相同 key 的请求结果同步阻塞版 * * param key 请求的唯一标识如 user_id、product_id、SQL 语句 * param loader 实际加载数据的函数如 DB 查询、RPC 调用 * return 加载结果 * throws Exception 加载过程中可能抛出的异常 * * 执行流程 * 1. 线程 A第一个请求调用 goFlight(user_123, loader) * → computeIfAbsent 发现没有 key创建一个新的 CompletableFuture * → 立即提交 loader 任务到线程池异步执行 * → 线程 A 调用 future.get() 阻塞等待结果 * * 2. 线程 B第二个请求几乎同时调用 goFlight(user_123, loader) * → computeIfAbsent 发现 key 已存在就是线程 A 创建的那个 Future * → 不会再执行 loader直接复用同一个 Future * → 线程 B 也调用 future.get() 阻塞等待结果 * * 3. loader 执行完成成功或失败触发 whenComplete 回调 * → 从 Map 中移除 key不管成功还是失败都要移除 * → 释放内存避免内存泄漏 * * 4. 线程 A 和线程 B 同时从 future.get() 醒来拿到结果 * * 关键设计点 * - computeIfAbsent 是原子操作保证了“检查-创建”的线程安全性 * - whenComplete 确保无论成功还是失败Map 都会被清理 * - future.get() 会阻塞当前线程但只有等待者会阻塞 * 真正执行 loader 的线程是线程池中的线程 */ public T goFlight(String key, FunctionString, T loader) throws Exception { // 步骤 1原子性地获取或创建 Future CompletableFutureT future ongoingRequests.computeIfAbsent(key, k - { // 这个 lambda 只在 key 不存在时执行由 computeIfAbsent 保证 // 返回一个新的 CompletableFuture并立即启动 loader 任务 // loader.apply(k) 会在 executor 线程池中异步执行 return CompletableFuture.supplyAsync(() - loader.apply(k), executor) .whenComplete((result, ex) - { // 步骤 2清理操作关键 // 无论 loader 执行成功还是抛出异常都要从 Map 中移除 key // 如果不移除后续所有请求都会无限等待一个已经完成的 Future // 但这里有个小问题如果有新的请求在“移除”和“完成”之间进来 // 可能会看到这个 key 还在 Map 中然后 get() 到一个已经完成的 Future // 这其实没问题因为 CompletableFuture 允许多次调用 get() // 只不过失去了“合并请求”的意义——但这也是可接受的折中。 // // 注意remove 时要确保移除的是当前这个 Future // 因为 Map 中的 key 可能已经被其他请求替换了 // 其实不会因为 computeIfAbsent 保证同一个 key 只会创建一个 Future // 直到它被移除后新的请求才会创建新的 Future。 ongoingRequests.remove(key); }); } ); // 步骤 3阻塞等待结果 // get() 会阻塞当前线程直到 CompletableFuture 完成 // 如果 loader 执行失败get() 会抛出 ExecutionException包装了原始异常 return future.get(); } // 异步版非阻塞 /** * 异步版本直接返回 CompletableFuture调用方自行决定如何等待 * * 适用场景调用方不想阻塞当前线程而是通过 thenApply / thenAccept 链式处理结果 * * 与 goFlight 的区别 * - goFlight阻塞式调用方线程会停在 future.get() 直到结果返回 * - goFlightAsync非阻塞式调用方拿到 Future 后自己编排后续逻辑 * * 设计考量的权衡 * - 阻塞式goFlight适合简单场景代码直观 * - 非阻塞式goFlightAsync适合响应式编程能更好地利用异步 I/O */ public CompletableFutureT goFlightAsync(String key, FunctionString, T loader) { // 核心逻辑与 goFlight 相同只是不调用 get()直接返回 Future // 这样调用方可以用 CompletableFuture 的链式 API 来处理结果 return ongoingRequests.computeIfAbsent(key, k - CompletableFuture.supplyAsync(() - loader.apply(k), executor) .whenComplete((result, ex) - { ongoingRequests.remove(key); }) ); } // 扩展方法 /** * 手动提前结束某个 key 的等待如超时清理 * * 注意调用此方法后所有等待该 key 的请求会收到一个 CancellationException。 * 但这只是“等待”被取消正在执行的 loader 可能还在后台运行无法强制中断。 * 更完善的方案需要引入 Context 或线程中断机制。 */ public void cancel(String key) { CompletableFutureT future ongoingRequests.remove(key); if (future ! null) { future.cancel(true); // true 表示允许中断正在执行的线程 } } /** * 获取当前正在处理的 key 数量用于监控 */ public int getOngoingCount() { return ongoingRequests.size(); } }Q1为什么要用CompletableFuture不能直接存结果吗如果直接存结果第 2 个人来的时候书还没找到你没法给他一个“等待”的机制。凭证CompletableFuture的好处是书没找到时凭证上写着“待完成”人可以等着。书找到后凭证自动变成“已完成”所有拿着凭证的人同时被唤醒。Q2为什么要用ConcurrentHashMap10 个人同时冲过来可能多人同时查本子。ConcurrentHashMap保证即使 10 个人同时查“这本书有没有人在找”也只会创建一张凭证不会创建 10 张。Q3为什么要异步执行loader仓库小哥找书如果前台管理员自己跑去仓库找书那后面 9 个人就只能干等着前台完全瘫痪。把找书交给仓库小哥线程池前台管理员可以继续接待其他请求系统不会被阻塞。使用示例public class Demo { public static void main(String[] args) { SingleFlightString sg new SingleFlight(); // 10 个并发请求同一个 key for (int i 0; i 10; i) { new Thread(() - { try { String result sg.goFlight(user_123, key - { // 这段代码只会执行一次 System.out.println(实际查询数据库...); Thread.sleep(100); return 用户数据_ System.currentTimeMillis(); }); System.out.println(结果: result); } catch (Exception e) { e.printStackTrace(); } }).start(); } } }五、实战Spring Boot 中防止缓存击穿5.1 完整实现import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cache.annotation.Cacheable; import org.springframework.stereotype.Service; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; Service public class UserService { // SingleFlight 实例可以注入为 Bean private final SingleFlightUser singleFlight new SingleFlight(); Autowired private UserRepository userRepository; /** * 查询用户信息 - 防缓存击穿版本 */ public User getUser(String userId) throws Exception { // 1. 先查缓存假设使用 Redis 或 Caffeine User cached getFromCache(userId); if (cached ! null) { return cached; } // 2. 缓存未命中使用 SingleFlight 保护 return singleFlight.goFlight(user: userId, key - { // 只有第一个请求会真正执行这里 System.out.println(从数据库加载用户: userId); User user userRepository.findById(userId) .orElseThrow(() - new RuntimeException(用户不存在)); // 加载完成后写入缓存 putToCache(userId, user); return user; }); } }5.2 配置为 Spring Beanimport org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; Configuration public class SingleFlightConfig { Bean public SingleFlightUser userSingleFlight() { // 使用自定义线程池避免影响其他业务 ExecutorService executor Executors.newFixedThreadPool(10); return new SingleFlight(executor); } }六、SingleFlight vs 其他方案方案并发请求处理下游压力适用场景无保护全部穿透N 倍❌ 不推荐互斥锁串行等待1 次低并发场景分布式锁串行等待 网络开销1 次分布式环境SingleFlight并发等待1 次✅ 高并发读场景布隆过滤器过滤不存在 key0 次防缓存穿透非击穿关键差异互斥锁是串行化请求排队SingleFlight 是并发等待所有请求同时等待同一个结果。在 1000 个并发请求的场景下互斥锁让 999 个请求排队等待而 SingleFlight 让 999 个请求同时被唤醒——响应时间天差地别。七、最佳实践与注意事项✅ 推荐做法合理设计 keykey 应该能唯一标识请求同时避免过于宽泛导致不该合并的请求被合并设置超时避免某个慢请求拖垮所有等待的线程使用自定义线程池避免默认线程池被阻塞任务耗尽监控dups指标观察请求合并效果评估是否需要调整配合缓存使用SingleFlight 是查不到缓存时的兜底结果应写入缓存八、总结SingleFlight 是一个小而美的并发控制模式核心价值将 N 个并发请求合并为 1 次实际调用大幅降低下游压力实现精髓Map WaitGroup Mutex的三重奏简洁而高效典型场景缓存击穿防护、接口防重、分布式锁优化