当前位置: 首页 > news >正文

SingleFlight 深度解析:从原理到源码,一文读懂请求合并利器

一、引言:那个让数据库瑟瑟发抖的瞬间

想象这样一个场景:凌晨三点,一个热点缓存 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 飙升

SingleFlight

1 次

并发等待,共享结果

关键在于:并发量不变,但下游压力降为 1

四、Java 实现 SingleFlight

4.1 基于 ConcurrentHashMap + CompletableFuture 的实现

这是最经典的 Java 实现方式,使用ConcurrentHashMap管理进行中的请求,CompletableFuture实现异步等待:

/** * SingleFlight - 请求合并器 * * 核心思想:当多个并发请求使用相同的 key 时,只让第一个请求真正执行, * 后续请求直接复用第一个请求的结果,从而避免重复的昂贵操作(如 DB 查询、RPC 调用)。 * * 类比:10 个人同时想查同一本书,只让 1 个人去仓库取,剩下 9 个人在原地等, * 等书取回来大家一起看。 * * @param <T> 返回结果的类型 */ public class SingleFlight<T> { // ==================== 核心存储结构 ==================== /** * 存储正在进行的请求:key → CompletableFuture<T> * * 为什么用 ConcurrentHashMap? * 因为会有多个线程同时尝试 put key,需要线程安全的 Map。 * * 为什么 key 对应的是 CompletableFuture,而不是直接存结果? * 因为后续请求需要能够“等待”结果,而不仅仅是拿到最终值。 * CompletableFuture 天然支持:多个线程可以同时调用 future.get(), * 一旦 future 完成(成功或失败),所有等待线程都会同时被唤醒。 */ private final ConcurrentHashMap<String, CompletableFuture<T>> 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, Function<String, T> loader) throws Exception { // ===== 步骤 1:原子性地获取或创建 Future ===== CompletableFuture<T> 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 CompletableFuture<T> goFlightAsync(String key, Function<String, 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) { CompletableFuture<T> future = ongoingRequests.remove(key); if (future != null) { future.cancel(true); // true 表示允许中断正在执行的线程 } } /** * 获取当前正在处理的 key 数量(用于监控) */ public int getOngoingCount() { return ongoingRequests.size(); } }

Q1:为什么要用CompletableFuture,不能直接存结果吗?

如果直接存结果,第 2 个人来的时候书还没找到,你没法给他一个“等待”的机制。凭证(CompletableFuture)的好处是:

  • 书没找到时,凭证上写着“待完成”,人可以等着。

  • 书找到后,凭证自动变成“已完成”,所有拿着凭证的人同时被唤醒。

Q2:为什么要用ConcurrentHashMap

10 个人同时冲过来,可能多人同时查本子。ConcurrentHashMap保证:即使 10 个人同时查“这本书有没有人在找”,也只会创建一张凭证,不会创建 10 张。

Q3:为什么要异步执行loader(仓库小哥找书)?

如果前台管理员自己跑去仓库找书,那后面 9 个人就只能干等着,前台完全瘫痪。把找书交给仓库小哥(线程池),前台管理员可以继续接待其他请求,系统不会被阻塞。

使用示例

public class Demo { public static void main(String[] args) { SingleFlight<String> 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 SingleFlight<User> 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 Bean

import 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 SingleFlight<User> userSingleFlight() { // 使用自定义线程池,避免影响其他业务 ExecutorService executor = Executors.newFixedThreadPool(10); return new SingleFlight<>(executor); } }

六、SingleFlight vs 其他方案

方案

并发请求处理

下游压力

适用场景

无保护

全部穿透

N 倍

❌ 不推荐

互斥锁

串行等待

1 次

低并发场景

分布式锁

串行等待 + 网络开销

1 次

分布式环境

SingleFlight

并发等待

1 次

✅ 高并发读场景

布隆过滤器

过滤不存在 key

0 次

防缓存穿透(非击穿)

关键差异:互斥锁是串行化(请求排队),SingleFlight 是并发等待(所有请求同时等待同一个结果)。在 1000 个并发请求的场景下,互斥锁让 999 个请求排队等待,而 SingleFlight 让 999 个请求同时被唤醒——响应时间天差地别。

七、最佳实践与注意事项

✅ 推荐做法

  1. 合理设计 key:key 应该能唯一标识请求,同时避免过于宽泛导致不该合并的请求被合并

  2. 设置超时:避免某个慢请求拖垮所有等待的线程

  3. 使用自定义线程池:避免默认线程池被阻塞任务耗尽

  4. 监控dups指标:观察请求合并效果,评估是否需要调整

  5. 配合缓存使用:SingleFlight 是"查不到缓存时的兜底",结果应写入缓存

八、总结

SingleFlight 是一个小而美的并发控制模式:

  • 核心价值:将 N 个并发请求合并为 1 次实际调用,大幅降低下游压力

  • 实现精髓Map + WaitGroup + Mutex的三重奏,简洁而高效

  • 典型场景:缓存击穿防护、接口防重、分布式锁优化

http://www.jsqmd.com/news/1304763/

相关文章:

  • 开源项目吐槽大会:一万字深度复盘,我亲历的那些“破防“瞬间
  • Android无线调试全链路排错指南:从ADB原理到实战解决连接问题
  • Zotero-OCR完全指南:免费PDF文字识别插件的终极解决方案
  • AI大模型的入门笔记
  • Python音频智能拼接:从基础概念到音频剧自动生成实战
  • PHP伪协议深度解析:从流操作到安全防御的完整指南
  • 电感核心公式V=L*(di/dt)深度解析与工程选型实战
  • 导师放养,硕士第一篇论文到底该从哪开始?
  • SIFT特征提取算法:原理、实现与OpenCV实战指南
  • LVDS接口技术解析:从差分信号原理到硬件设计实战
  • 2026 年 7 月新发布:京口优秀的透光混凝土板 厂商哪家可靠,如果把建筑墙面上的这块透光“薄石板”,换成能把自然光带进地下室的好东西?-石美清水混凝土板 - 企业推荐管【认证】
  • Ansible Playbook核心概念与高级特性实战指南
  • 【智能体安全治理|专栏第8期】:智能体安全攻防全景图:我们实践中的六层攻击面与真实对抗经验
  • 告别激光扫描与人工建模:无前置建模动态三维重构的算力革命研发课题方案
  • 2026年毕业生黑科技榜单9款AI论文平台横评!
  • Unity时间系统深度解析:从Time.deltaTime到自定义时间层
  • Maya角色绑定、蒙皮与权重调整全流程实战指南
  • C++ vector::erase迭代器失效与安全删除模式详解
  • 基于LLM与语音技术的智能客服系统构建实战
  • 三菱FX2N-2DA模拟量输出模块:从硬件接线到编程调试的完整指南
  • CAN总线ESD保护设计实战:从TVS选型到PCB布局的避坑指南
  • 从水管网络到算法实现:深入理解最大流与最小割的核心原理与应用
  • Linux 终端快捷键
  • 从Python到CUDA,AI文件读写的7层加速架构,含TensorFlow/PyTorch原生适配清单(限首批开源)
  • 2026年陕西电力建筑新能源企业咨询服务推荐榜:资质升级规划、人员补充、证照维护、注册人员转入转出及方案设计政策咨询优选 - 优企名品
  • Vin象棋:三分钟搭建你的AI象棋助手,免费体验专业级对弈指导
  • 一次消息如何变成多步行动:拆开 OpenClaw 的 Agent Loop
  • JasperReports报表引擎实战:从模板设计到Spring Boot集成与性能优化
  • OMG数据集:多模态基因组语言模型的数据基石与实战指南
  • Qt QSS样式表开发实战与性能优化指南