系列目录
- 从单线程到线程池:云盘转 Java 后的第一堂并发课
- 线程池不是 new 出来就完事:参数、队列与快慢接口隔离
- 数据库连接池与 Spring 声明式事务:把 Node.js 的坑填上
- ConcurrentHashMap 与锁:文件元数据的并发读写
- Future 与 CountDownLatch:一个接口聚合一堆下游(本篇)
先补齐聚合调用的知识
从 Node.js、C# 等后端经验转到 Java 后,我预期会遇到两类问题:共享资源在多线程下的竞争,以及一个接口依赖多个下游后的等待时间。前几篇已经处理了线程池、连接池和共享数据的访问方式。聚合接口的下游数量会继续增加,因此在项目需要改造前,先学习了 Future、Callable 和 CountDownLatch,并用首页接口验证它们的适用范围。
云盘客户端首页由 web 端和 PC 端共用。接口需要返回默认目录的文件列表、用户配额和回收站占用统计,后来又增加了「最近动态摘要」。每一项都对应一次下游调用。列表和元数据走自己的 service,配额和回收站走另一个模块,动态摘要还要调用消息侧服务。
第一版按顺序调用,拿到结果后写入返回对象:
public HomeView getHome(String uid) {
HomeView view = new HomeView();
view.setFileList(fileService.listDefaultDir(uid)); // ~80ms
view.setQuota(quotaService.getQuota(uid)); // ~150ms
view.setRecycleStat(recycleService.stat(uid)); // ~60ms
return view;
}接口刚上线时,下游较少且响应较快,RT 是两三百毫秒。下游增加后,RT 也随之增加。加第四个下游后,首页 P95 已接近一秒。日常均值只有五六百毫秒,但下游抖动会拉高尾部延迟,产品同学开始反馈首页变慢。
当时先和产品核对返回字段,删掉了一个「有更好、没有也行」的非必要字段,将下游从四个减为三个,RT 回到可接受范围。这个处理只能减少当前的串行等待时间。串行调用的耗时近似为各下游耗时之和,新增一个依赖通常会增加接口的等待时间。
提前准备的 demo 用来对比串行和并发聚合,JDK 7 可以运行:
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
public class AggregateDemo {
// 模拟一次下游调用:sleep 指定毫秒数后返回结果
static String callDownstream(String name, long costMs) {
try {
Thread.sleep(costMs);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
return name + " 的结果";
}
public static void main(String[] args) throws Exception {
// 下游清单:名字 + 各自耗时
final String[] names = {"文件列表", "配额", "回收站统计", "最近动态"};
final long[] costs = {80, 150, 60, 200};
// 方式一:串行
long t1 = System.currentTimeMillis();
List<String> r1 = new ArrayList<String>();
for (int i = 0; i < names.length; i++) {
r1.add(callDownstream(names[i], costs[i]));
}
System.out.println("串行耗时:" + (System.currentTimeMillis() - t1) + "ms");
// 方式二:Future 并发聚合
ExecutorService pool = Executors.newFixedThreadPool(4);
long t2 = System.currentTimeMillis();
List<Future<String>> futures = new ArrayList<Future<String>>();
for (int i = 0; i < names.length; i++) {
final String name = names[i];
final long cost = costs[i];
futures.add(pool.submit(new Callable<String>() {
@Override
public String call() {
return callDownstream(name, cost);
}
}));
}
List<String> r2 = new ArrayList<String>();
for (Future<String> f : futures) {
r2.add(f.get(500, TimeUnit.MILLISECONDS)); // 必须带超时,后面讲
}
System.out.println("并发耗时:" + (System.currentTimeMillis() - t2) + "ms");
pool.shutdown();
}
}这个 demo 中,串行耗时约为 490ms,即 80、150、60 和 200 的总和。线程池有 4 个线程,4 个任务都能立即开始,并发耗时通常略高于 200ms,接近最慢任务的耗时。实际接口还会受到线程池排队、网络调度和序列化等因素影响。并发能重叠独立下游的等待时间,慢下游的服务时间仍然存在。
Callable 与 Future 如何拆分提交和取结果
第 2 篇中的线程池任务主要使用 Runnable。Runnable.run() 没有返回值,也不能声明受检异常。聚合场景需要将下游结果写入响应对象,因此使用 Callable 和 Future 更合适。
Callable表示可执行任务,call()有返回值,并且可以抛出异常。Future是ExecutorService.submit(Callable)返回的句柄。提交成功后,任务可能正在执行、排队、已完成或已取消。submit通常很快返回,future.get()在结果尚未完成时等待,并在完成后返回结果或抛出异常。
聚合分为两段。先提交全部任务,再读取结果。第一个循环结束后,任务已进入线程池;第二个循环等待某个结果时,其他任务仍可继续执行。任务能并发开始且相互独立时,总等待时间接近最慢任务的耗时。
任务提交与读取结果的顺序会影响并发度。每次 submit 后立即 get 时,调用线程先等待当前任务完成,后续任务尚未提交,执行效果接近串行。代码 review 时需要检查这两个操作是否写在同一个循环中。
CountDownLatch 与 Future.get 分别传递什么信息
CountDownLatch 经常和 Future 一起出现,但两者传递的信息不同。
Future.get用于取得一个任务的完成状态和结果。任务异常时,调用方也能通过ExecutionException感知失败。CountDownLatch只表示计数是否归零。等待线程通过await()等待 N 个任务调用countDown(),它不保存任务结果,也不传递任务异常。
聚合接口通常需要每个下游的结果,Future 更直接。CountDownLatch 适合只等待一组工作完成的场景,例如批量处理文件后由主线程统一生成处理报告。报告需要逐项结果时,任务还需要把结果和异常写入线程安全的容器。CountDownLatch 不能重置,同一个实例只能使用一次。
await() 和 get() 都提供无限等待与带超时的形式。请求链路需要明确等待上限,避免下游未返回时长期占用请求线程。
get 的超时、总预算与取消边界
Future 有不带参数的 get(),也有 get(timeout, unit)。后者在指定时间内未完成时抛出 TimeoutException。需要在响应时间预算内返回的聚合接口,应使用有界等待,并处理超时、取消、执行异常和中断。
并发后,独立下游的等待时间可以重叠,但接口仍要等待所需结果或相应的超时。最慢下游会影响整体 RT。下游服务 hang 住且没有超时时,调用不带参数的 get() 会持续等待,请求线程也会持续被占用。
每个下游可以配置超时预算。超时后记录日志,并将对应字段降级为空或默认值,其余字段继续返回。这里有两个预算。单次 Future.get(timeout, unit) 限制一次等待;按顺序对多个 Future 分别等待时,整个聚合等待时间可能超过任一单独超时。线上实现还应按接口总 deadline 计算每次剩余等待时间,或使用统一的超时协调方式。
get 超时不会自动停止对应任务,任务仍可能在线程池中执行。可以调用 future.cancel(true) 请求取消:尚未开始的任务可能不会执行,正在执行的任务会收到中断请求。任务和底层客户端是否检查并正确处理中断,决定了任务能否停止,cancel 也可能因任务已经完成或已取消而返回 false。
我们的 HTTP/RPC 客户端大多对中断不敏感,取消不能替代下游调用自身的连接、读取或 RPC 超时。下游没有 socket 超时时,Future 层的超时只能让等待方先返回,底层连接和执行线程仍可能继续占用资源。这个限制在系列收尾讨论超时体系时再展开。
聚合任务不能与父任务竞争同一个业务池
并发化改造的第一版为了省事,将聚合子任务提交到业务线程池。自测正常,RT 也降了下来。预发压测时,首页流量上升后,其他接口开始排队。
当时的结构是:Jetty 线程接收请求后,将整个请求处理,包括聚合逻辑,提交给业务池。处理首页请求的业务线程又向同一个池提交 4 个子任务,然后等待结果。子任务和普通请求共享线程与队列,首页请求会增加池内任务数;等待中的父任务还会占住业务线程。高并发下,池中的工作线程都被等待子任务的父任务占住时,子任务只能留在队列中,父任务也无法继续执行。这是线程饥饿死锁的一种典型形态。带超时的 get 会限制一次等待,线程池结构造成的竞争仍然存在。
改法是为聚合任务使用独立线程池,与业务入口池隔离。第 2 篇提到过,CallerRunsPolicy 只用于内部任务池。对于配置了有界队列的聚合池,线程数达到上限且队列也满时,CallerRunsPolicy 会让提交任务的业务线程同步执行任务。提交速度会降低,任务不会继续被拒绝或无限进入队列。这会增加该请求的耗时,也会占用业务线程,是否合适取决于接口的超时预算和业务入口池的余量。业务入口池仍使用 AbortPolicy 快速失败。
隔离后,聚合池的容量可以单独评估。一个粗略的并发需求估算是:聚合接口峰值 QPS × 每请求提交的下游任务数 × 下游平均服务时间。它是 Little 定律在稳定平均条件下的近似,实际配置还要考虑目标延迟、P99、队列长度、拒绝策略和下游限流,不能只按这个乘积确定线程数。
项目中的改造与监控
云盘最终做了三项改造。
- 聚合接口并发化。首页和另外两处拼多个下游的接口改为先提交全部任务、再读取结果,子任务使用独立的聚合线程池。首页 RT 从接近一秒回到两百多毫秒。后续增加字段时,接口等待时间还取决于新下游的耗时、线程池排队和总超时预算,不再简单按所有下游耗时相加。
- 每个下游一个超时配置。初始值参考各自的历史 P99。超时字段降级返回,列表类字段给空结果,统计类字段给默认值。字段级降级日志单独记录,便于定位超时下游。
- 监控聚合池。活跃线程、队列长度和拒绝次数接入第 2 篇使用的同一套公司监控平台,也记录
CallerRunsPolicy的触发次数。该指标上升说明池已发生饱和,需要结合下游耗时、流量和队列情况判断是扩容、限流还是减少聚合任务。
线程池隔离减少接口之间的资源竞争,Future 的有界等待限制单个接口受某个下游影响的时间。它们处理的资源占用范围不同。
当时只在 demo 中使用 CompletableFuture
CompletableFuture 是 JDK 8 的 API。2016 年 JDK 8 已发布两年,线上仍运行 JDK 7。为提前了解后续可用的组合方式,我在本地 demo 中试过 thenCombine 和 allOf 等链式组合。它们可以表达任务依赖和结果汇总,减少手写 Future 循环。
JDK 8 的 CompletableFuture 没有 orTimeout 或 completeOnTimeout 这类内置超时 API。JDK 8 中需要通过额外的定时任务或自行协调实现超时和降级;这些 API 是 JDK 9 才提供的。因此,当时的 demo 没有改变线上超时方案。
线上是 JDK 7。为一个聚合接口升级 JDK 不在当时的排期内。手写 Future 的版本满足需求,团队也能维护。CompletableFuture 后来在另一个项目中使用过,这个系列不再展开。
参考资料
- JDK 7
Future/ExecutorService/CountDownLatchJavadoc - 《Java并发编程实战》(Brian Goetz)第 6 章

