Skip to content
Charles Shao
Go back

并发协作与异步编排:闭锁、信号量与 CompletableFuture

views

有了锁和并发容器,还差一块:线程之间怎么协同——等一批任务完成、多线程在某点集合再一起走、限制同时访问某资源的线程数;以及在异步世界里,怎么把「调 A 再调 B、并行调 C 和 D 再合并」串成流水线。

前半讲三个协作工具(CountDownLatch / CyclicBarrier / Semaphore),后半讲 CompletableFuture——广告竞价里「并行 fan-out 到多个下游、带超时地聚合」的标准写法。

TL;DR

Table of contents

Open Table of contents

1. 三个协作工具

三个协作工具对比。CountDownLatch:等计数归零、一次性、AQS 共享;CyclicBarrier:N 线程互等到齐、可循环、Lock+Condition;Semaphore:许可数限流、可重用、AQS 共享。底部说明选型与底层差异。

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=计数
CyclicBarrierN 线程互等到齐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); // 转换结果(异步,指定线程池)

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());

3. 实战:竞价并行拉取多路依赖

竞价要在 tmax(如 100ms)内并行拉用户画像、定向特征、频控状态,任何一路超时都降级,最后聚合出价:

竞价 fan-out 示意图。左侧 BidRequest 主线程扇出到三路 supplyAsync:画像(超时→empty)、特征(超时→defaults)、频控(超时→放行);侧注 allOf 聚合、任一路可降级、总耗时约等于最慢一路。底部强调隔离线程池与双保险超时。

// 关键:为不同下游用隔离的线程池(舱壁),别用默认 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() 抛非受检的 CompletionExceptionget() 抛受检的 ExecutionException/InterruptedException。在 lambda/流里通常用 join() 更顺手,但要记得它仍是阻塞调用——别在异步回调里对另一个未完成的 CF join(),容易把线程池线程占死甚至死锁。

JMM 的可见性论证,到锁的两条路线、线程池并发容器,再到本篇的协作与异步编排——它们共同支撑高并发服务在有限线程与资源下的正确性与吞吐。继续深入运行时,可看 连接与 I/O 模型 与 Netty。

延伸阅读


views
Share this post on:

Previous Post
Netty 深挖(开篇)· Reactor 模型与 EventLoop 线程模型
Next Post
并发容器:ConcurrentHashMap 与阻塞队列