有了锁和并发容器,还差一块:线程之间怎么协同——等一批任务完成、多线程在某点集合再一起走、限制同时访问某资源的线程数;以及在异步世界里,怎么把「调 A 再调 B、并行调 C 和 D 再合并」串成流水线。
前半讲三个协作工具(CountDownLatch / CyclicBarrier / Semaphore),后半讲 CompletableFuture——广告竞价里「并行 fan-out 到多个下游、带超时地聚合」的标准写法。
TL;DR
- CountDownLatch:一次性闭锁,
await()等到计数被countDown()减到 0。用完即废,不能重置。典型:主线程等所有子任务完成。(AQS 共享) - CyclicBarrier:可循环的栅栏,N 个线程互等到齐再一起继续,可重复使用。(
ReentrantLock+Condition,不是 AQS) - Semaphore:管理 N 个许可,用于限制并发数。(AQS 共享)
- Future 的局限:
get()只能阻塞,难以组合与回调。CompletableFuture 补齐链式编排、组合、异常与超时。 - 默认线程池是
ForkJoinPool.commonPool(),做阻塞 I/O 会拖垮全局——务必传自定义线程池,并按业务隔离。 - 超时兜底(
orTimeout/completeOnTimeout,Java 9+)让异步调用不会无限挂起。
Table of contents
Open Table of contents
1. 三个协作工具
1.1 CountDownLatch:等一批任务完成
计数器初始化为 N,每个任务完成调 countDown() 减一,等待方 await() 阻塞到计数归零。一次性,不可重置。
int taskCount = 5;
CountDownLatch latch = new CountDownLatch(taskCount);
ExecutorService pool = new ThreadPoolExecutor(
5, 5, 0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(taskCount),
new ThreadPoolExecutor.CallerRunsPolicy());
for (int i = 0; i < taskCount; i++) {
final int id = i;
pool.execute(() -> {
try {
doSubTask(id);
} finally {
latch.countDown(); // 无论成功失败都要减,放 finally
}
});
}
latch.await(3, TimeUnit.SECONDS); // 主线程等全部完成(带超时,避免死等)
场景:主流程需要等多个并行子任务全部返回再汇总;或”发令枪”用法——N 个线程先 await() 一个初值为 1 的 latch,主线程 countDown() 让它们同时起跑(压测常用)。
1.2 CyclicBarrier:多线程互相等待,可循环
N 个线程各自 await(),都到齐后一起放行,然后栅栏自动重置可再用。还能传一个 barrierAction,在最后一个线程到达时(放行前)执行一次汇总:
int parties = 3;
CyclicBarrier barrier = new CyclicBarrier(parties, () ->
System.out.println("本轮 3 个线程都到齐,执行汇总")); // 到齐后的回调
for (int i = 0; i < parties; i++) {
new Thread(() -> {
for (int round = 0; round < 2; round++) { // 可循环两轮
computePhase(round);
try {
barrier.await(); // 等其他线程,到齐一起进入下一轮
} catch (Exception e) { Thread.currentThread().interrupt(); }
}
}).start();
}
CountDownLatch vs CyclicBarrier:前者是”一个/一批线程等别的任务完成”,计数减到 0、一次性;后者是”一组线程互相等对方”,到齐即放行、可复用。分阶段迭代计算(如并行分块后每轮同步一次)用 CyclicBarrier。
1.3 Semaphore:限制并发数
信号量维护 N 个许可,acquire() 取一个(无则阻塞),release() 还一个。用来给并发访问设上限——限流、连接/资源池:
Semaphore semaphore = new Semaphore(10); // 最多 10 个线程同时访问下游
public Result callDownstream(Request req) throws InterruptedException {
if (!semaphore.tryAcquire(50, TimeUnit.MILLISECONDS)) {
return Result.rejected(); // 拿不到许可 → 快速失败/降级
}
try {
return downstream.call(req); // 受保护的临界资源
} finally {
semaphore.release(); // 必须在 finally 归还
}
}
场景:保护脆弱下游不被打爆(一种轻量限流,和 服务限流 的令牌桶互补——Semaphore 控”并发数”,令牌桶控”速率”);实现简单的对象/连接池。
三者速查:
| 工具 | 语义 | 可重用 | 底层 |
|---|---|---|---|
| CountDownLatch | 等计数归零 | ❌ 一次性 | AQS 共享,state=计数 |
| CyclicBarrier | N 线程互等到齐 | ✅ | ReentrantLock + Condition |
| Semaphore | 控制并发许可数 | ✅ | AQS 共享,state=许可 |
2. 从 Future 到 CompletableFuture
Future(JDK 5)能拿到异步结果,但能力很窄:
Future<Profile> f = pool.submit(() -> loadProfile(uid));
Profile p = f.get(); // 只能阻塞等待;想"完成后自动做下一步"?做不到
它的痛点:get() 只能阻塞;无法把多个 Future 组合(“A 完成后用结果调 B”);无法注册回调;异常处理笨拙(要 try-catch 包 get)。
CompletableFuture(JDK 8)实现了 Future + CompletionStage,把异步编程变成声明式的流水线:任务完成会自动触发下一环,全程不阻塞。
2.1 创建与串行编排
CompletableFuture<Profile> cf =
CompletableFuture.supplyAsync(() -> loadProfile(uid), pool) // 异步执行,有返回值
.thenApply(p -> enrich(p)) // 转换结果(同步,复用上一步线程)
.thenApplyAsync(p -> score(p), pool); // 转换结果(异步,指定线程池)
supplyAsync(有返回值)/runAsync(无返回值):提交任务。thenApply/thenAccept/thenRun:拿上一步结果做转换 / 消费 / 收尾。- 带
Async后缀的版本在另一个线程(可指定线程池)执行,不带的复用上一环的线程。
2.2 依赖与组合
// thenCompose:A 的结果决定 B(扁平化,避免 CF<CF<T>> 嵌套)
CompletableFuture<Order> order =
loadUserAsync(uid).thenCompose(user -> loadOrderAsync(user.id()));
// thenCombine:A、B 并行,都完成后合并
CompletableFuture<Page> page =
loadUserAsync(uid).thenCombine(loadAdsAsync(uid), (user, ads) -> render(user, ads));
thenCompose ≈ flatMap(串行依赖),thenCombine ≈ 两路并行后 zip。
2.3 聚合多个:allOf / anyOf
CompletableFuture<Resp> a = callDspA(req);
CompletableFuture<Resp> b = callDspB(req);
CompletableFuture<Resp> c = callDspC(req);
// 等全部完成
CompletableFuture<Void> all = CompletableFuture.allOf(a, b, c);
all.thenRun(() -> {
List<Resp> results = Stream.of(a, b, c)
.map(CompletableFuture::join).collect(toList()); // 已完成,join 不阻塞
pickWinner(results);
});
// 任意一个完成即返回(如多副本取最快)
CompletableFuture<Object> fastest = CompletableFuture.anyOf(a, b, c);
2.4 异常与超时
CompletableFuture<Resp> resp = callDsp(req)
.orTimeout(20, TimeUnit.MILLISECONDS) // JDK 9+:超时则以 TimeoutException 完成
.exceptionally(ex -> Resp.noBid()) // 任何异常 → 降级为 no-bid
// 或用 handle 同时处理正常/异常两种结果:
.handle((r, ex) -> ex == null ? r : Resp.noBid());
exceptionally:只在异常时兜底。handle/whenComplete:同时拿到结果和异常(whenComplete不改变结果,handle可转换)。orTimeout(超时抛异常)/completeOnTimeout(超时给默认值):Java 9+,高并发尾延迟控制的关键。JDK 8 需自己用ScheduledExecutorService实现超时。
3. 实战:竞价并行拉取多路依赖
竞价要在 tmax(如 100ms)内并行拉用户画像、定向特征、频控状态,任何一路超时都降级,最后聚合出价:
// 关键:为不同下游用隔离的线程池(舱壁),别用默认 commonPool
private final ExecutorService featurePool = /* 见线程池篇 */ ...;
public BidResponse bid(BidRequest req) {
long tmax = 90; // 给聚合留余量,小于外部 tmax 100ms
CompletableFuture<Profile> profileF = CompletableFuture
.supplyAsync(() -> profileService.get(req.uid()), featurePool)
.completeOnTimeout(Profile.empty(), tmax, MILLISECONDS) // 超时给空画像
.exceptionally(ex -> Profile.empty()); // 异常也降级
CompletableFuture<Features> featF = CompletableFuture
.supplyAsync(() -> featureService.get(req), featurePool)
.completeOnTimeout(Features.defaults(), tmax, MILLISECONDS)
.exceptionally(ex -> Features.defaults());
CompletableFuture<Boolean> capF = CompletableFuture
.supplyAsync(() -> freqCapService.allow(req.uid()), featurePool)
.completeOnTimeout(true, tmax, MILLISECONDS) // 超时默认放行
.exceptionally(ex -> true);
// 三路并行,全部(降级)完成后聚合出价
return CompletableFuture.allOf(profileF, featF, capF)
.thenApply(v -> computeBid(profileF.join(), featF.join(), capF.join()))
.exceptionally(ex -> BidResponse.noBid())
.join(); // 竞价主线程在此同步等待最终结果
}
要点:① 三路 supplyAsync 并行发起,总耗时 ≈ 最慢的一路而非三路之和;② 每一路都带 completeOnTimeout + exceptionally 双保险,任何下游抖动都不拖垮整体,保住 tmax;③ 用隔离的业务线程池而非 commonPool;④ 只在最外层 join() 一次同步等待。这套”并行 fan-out + 超时降级 + 聚合”是高并发低延迟服务的标准骨架,与 依赖韧性 的超时/降级、连接与 I/O 模型 的阻塞卸载一脉相承。
4. CompletableFuture 的坑
4.1 默认线程池的陷阱
不带 Async 或不传 executor 的方法,用的是 ForkJoinPool.commonPool()——一个全 JVM 共享、线程数默认约等于”核数-1”的池。在里面跑阻塞 I/O 会:① 迅速占满,拖垮所有用它的地方(包括 parallel stream);② 线程数太少,I/O 密集任务吞吐上不去。
约定:只要任务涉及阻塞/I/O,一律传自定义线程池,并按业务隔离(见 线程池原理与调优 的舱壁):
CompletableFuture.supplyAsync(() -> blockingCall(), myIoPool); // ✅ 显式线程池
4.2 异常被吞
如果一条链最后没有 exceptionally/handle/whenComplete,异常会被封在返回的 CF 里,直到你 get/join 才抛——中途静默丢失,日志里什么都看不到。每条异步链都要有终点的异常处理。
4.3 join vs get
join() 抛非受检的 CompletionException,get() 抛受检的 ExecutionException/InterruptedException。在 lambda/流里通常用 join() 更顺手,但要记得它仍是阻塞调用——别在异步回调里对另一个未完成的 CF join(),容易把线程池线程占死甚至死锁。
从 JMM 的可见性论证,到锁的两条路线、线程池 与 并发容器,再到本篇的协作与异步编排——它们共同支撑高并发服务在有限线程与资源下的正确性与吞吐。继续深入运行时,可看 连接与 I/O 模型 与 Netty。
延伸阅读
- JMM 与可见性
- synchronized 与锁升级
- AQS 与 Lock 家族
- 线程池原理与调优
- 并发容器
- 连接与 I/O 模型
- OpenJDK 源码:
CompletableFuture、CountDownLatch、CyclicBarrier、Semaphore - Brian Goetz et al. Java Concurrency in Practice, Ch. 5 & 8
- Oracle. CompletionStage / CompletableFuture API 文档