CompletableFuture 并发统计
一、问题:串行统计太慢
很多后台系统都有"工作台/仪表盘"页面,需要一次性展示多个统计数字:在职人数、离职人数、证件到期人数、员工生日人数……
最朴素的写法是串行调用:
long dimission = service.getDimissionTotal(param); // 150ms long onJob = service.getCount(param); // 200ms long idCard = service.getEmployeeIdCardExpireCount(); // 180ms long birthday = service.getEmployeeBirthdayCount(); // 120ms问题:4 次数据库查询排队执行,总耗时 = 150 + 200 + 180 + 120 =650ms。
但这 4 个查询互相独立、没有依赖关系——完全可以同时跑。
二、并发:让 4 个查询同时执行
2.1 时间线对比
串行(650ms): 在职(200ms) → 离职(150ms) → 证件到期(180ms) → 生日(120ms) 并发(≈200ms): 在职(200ms) ┐ 离职(150ms) ┤ 证件到期(180ms)┼─ 同时跑,总耗时 ≈ 最慢的那个 生日(120ms) ┘并发后总耗时从 650ms 降到约200ms(最慢任务的耗时),提速 3 倍多。
2.2 代码实现
EmployeeTotalVo allTotal = new EmployeeTotalVo(); CompletableFuture.allOf( CompletableFuture.runAsync(() -> allTotal.setDimissionTotal( String.valueOf(service.getDimissionTotal(param)))), CompletableFuture.runAsync(() -> allTotal.setBeonTheJobTotal( String.valueOf(service.getCount(param)))), CompletableFuture.runAsync(() -> allTotal.setIdcardExpireTotal( service.getEmployeeIdCardExpireCount(param))), CompletableFuture.runAsync(() -> allTotal.setStaffBirthdayTotal( service.getEmployeeBirthdayCount(param))) ).join(); return RequestResult.success(allTotal);三、逐个 API 拆解
3.1runAsync—— 异步启动任务
CompletableFuture.runAsync(() -> 任务代码)- 立即返回一个
CompletableFuture,不阻塞当前线程; - 任务在独立线程(默认
ForkJoinPool.commonPool())里执行; - 4 行
runAsync= 启动 4 个任务,几乎同时开始。
💡 这里的
() -> ...是 Lambda,即"这个任务具体干什么"。
3.2allOf—— 组合多个任务
CompletableFuture.allOf(任务1, 任务2, 任务3, 任务4)把多个CompletableFuture打包成一个"组合任务",它代表**"这几个全都完成"**这个条件。
还有个孪生兄弟:
| API | 放行条件 |
|---|---|
allOf(...) | 全部完成 |
anyOf(...) | 任意一个完成即可 |
3.3join—— 阻塞等待
allOf(...).join();阻塞当前线程,直到组合任务完成(即 4 个子任务都跑完)。.join()之后的代码,必然是 4 个统计都已填入allTotal才会执行。
并发安全性说明:allTotal虽然是同一个对象,但 4 个任务设置的是不同字段(setDimissionTotal、setBeonTheJobTotal...),互不干扰,因此并发安全。
四、join的同类 API 全景
"等待异步任务完成"这一类,Java 提供了多种选择。理解它们的差异,才能在工程中选对。
4.1CompletableFuture家族
| API | 是否阻塞 | 异常处理 | 适用场景 |
|---|---|---|---|
join() | ✅ 阻塞 | 抛CompletionException(unchecked) | 业务代码首选,简洁 |
get() | ✅ 阻塞 | 抛 checked 异常(必须 try-catch) | 兼容老 API |
getNow(默认值) | ❌ 不阻塞 | 不阻塞 | 没完成就返回默认值 |
get(timeout, unit) | ✅ 阻塞(带超时) | 抛TimeoutException | 防止慢任务卡死 |
4.2joinvsget的核心区别
// join():不抛 checked 异常,代码干净 future.join(); // 直接用,无需 try-catch // get():被迫处理 checked 异常 try { future.get(); } catch (InterruptedException | ExecutionException e) { // 被迫写 try-catch,代码啰嗦 }结论:业务代码优先用join(),避免为了 checked 异常写一堆样板代码。
4.3 传统并发等待(底层/老代码)
| API | 特点 |
|---|---|
Future.get() | ExecutorService.submit()返回,功能类似但不能链式编排 |
CountDownLatch.await() | 计数器到 0 才放行,适合"等 N 个线程都完成" |
CyclicBarrier.await() | 所有线程到齐一起继续,可复用 |
Thread.join() | 等某个线程死亡,最底层 |
4.4 不阻塞的回调式编排
在响应式编程(WebFlux、Reactor)中,不能阻塞线程,改用回调链:
future .thenApply(result -> 加工结果) // 完成后转换 .thenCompose(result -> 另一个异步) // 完成后串联另一个异步 .thenAccept(result -> 消费结果); // 完成后消费这种方式不阻塞线程,吞吐量更高,但调试更复杂。
五、工程选型建议
| 场景 | 推荐 |
|---|---|
| 多个独立任务并发,等全部完成 | CompletableFuture.allOf(...).join()✅(本文场景) |
| 需要超时保护,防止慢任务卡死 | .get(timeout, unit)或orTimeout()(Java 9+) |
老项目用线程池submit | Future.get(),但建议迁移到CompletableFuture |
| 高并发响应式系统 | .thenApply()/.thenCompose()回调编排 |
| "等 N 个线程都完成"的通用同步 | CountDownLatch |
六、进阶:给统计接口加超时保护
生产环境有个隐患:如果某个 count 查询卡住(比如数据库锁、慢查询),.join()会无限等待,整个接口超时,前端一直转圈。
改进方案——加超时:
try { CompletableFuture.allOf(f1, f2, f3, f4) .get(3, TimeUnit.SECONDS); // 最多等 3 秒 } catch (TimeoutException e) { // 超时后,已完成的字段有值,未完成的保持默认(0) log.warn("工作台统计部分超时,返回已完成的字段"); } catch (InterruptedException | ExecutionException e) { log.error("工作台统计异常", e); } return RequestResult.success(allTotal);这样即使某个统计卡死,接口也能在 3 秒内返回(已完成的有值,未完成的用默认 0),保证可用性优于精确性。
七、可读性优化
原代码把 4 个任务塞在allOf参数里,稍显拥挤。可以拆开声明,逻辑完全等价但更易读:
CompletableFuture<Void> f1 = CompletableFuture.runAsync( () -> allTotal.setDimissionTotal(String.valueOf(service.getDimissionTotal(param)))); CompletableFuture<Void> f2 = CompletableFuture.runAsync( () -> allTotal.setBeonTheJobTotal(String.valueOf(service.getCount(param)))); CompletableFuture<Void> f3 = CompletableFuture.runAsync( () -> allTotal.setIdcardExpireTotal(service.getEmployeeIdCardExpireCount(param))); CompletableFuture<Void> f4 = CompletableFuture.runAsync( () -> allTotal.setStaffBirthdayTotal(service.getEmployeeBirthdayCount(param))); CompletableFuture.allOf(f1, f2, f3, f4).join();八、总结
| 知识点 | 要点 |
|---|---|
| 为什么要并发 | 独立任务并发执行,总耗时 ≈ 最慢的任务,而非累加 |
| 三件套 | runAsync(启动)→allOf(组合)→join(等待) |
joinvsget | join不抛 checked 异常,业务代码首选 |
allOfvsanyOf | 全部完成 vs 任意一个完成 |
| 生产加固 | 用.get(timeout)加超时,防止单点慢任务拖垮接口 |
| 线程安全前提 | 并发写同一对象的不同字段才安全;写同字段需加锁或用原子类 |
核心心智模型:runAsync负责"派活",allOf负责"收口",join负责"等结果"。三个组合起来,就是"派多个活、同时干、干完一起收"的并发统计范式。
这套模式非常适合"工作台/仪表盘/报表汇总"等需要聚合多个独立数据源的场景。关键是判断任务之间是否真的独立——有依赖关系就不能简单并发,要用
thenCompose串联。
