以并行查询和本地缓存为例,说明 CompletableFuture 如何选择线程池、为什么在线程池任务中调用 join 可能造成线程饥饿,以及 ReentrantLock 应该如何划定边界,最后给出可运行的组合实践。

问题背景:异步不等于不会阻塞

在 Java 服务中,CompletableFuture 经常被用来并行调用多个下游接口:用户信息、库存、价格、优惠券分别查询,最后合并成一个结果。代码看起来从同步调用变成了异步调用,但真正的问题通常出现在三个地方:

  1. 所有异步任务都落到同一个线程池,慢任务把快任务也拖住。
  2. 在线程池任务内部调用 join()get(),导致工作线程被占满。
  3. 持有锁时等待另一个异步阶段,而异步阶段又需要竞争同一把锁。

这几类问题往往不会在单元测试中立刻暴露。请求量上升、下游变慢或线程池缩小之后,服务才会出现响应时间突然升高,甚至所有任务都停在等待状态。

先明确三个边界

1. CompletableFuture 只负责描述流程,不负责自动创建合理的线程模型

runAsyncsupplyAsync 如果没有显式传入 Executor,会使用 JDK 提供的默认公共线程池。这个线程池适合较短的 CPU 型任务,不适合承载需要网络等待、数据库等待的业务操作。

工程代码中,通常应该至少区分两类执行器:

  • IO 线程池:用于 HTTP、数据库、文件等阻塞操作,线程数可以根据下游容量和等待时间评估。
  • CPU 线程池:用于计算、组装、序列化等短任务,线程数不宜盲目扩大。

线程池不是越多越好。线程池之间如果没有明确职责,只是把问题从一个池子复制成多个池子,最终仍然会因为下游过载或上下文切换增加而变慢。

2. 锁只保护共享状态,不保护整个异步流程

锁的职责应当是保护一小段对共享数据的读写。网络调用、数据库调用、join()、重试和复杂计算,都不应该放在锁的临界区内。

一个简单判断方式是:如果拿掉这把锁,代码中是否还需要等待一个外部操作完成?如果答案是肯定的,通常说明锁的范围过大。

3. join 的安全性取决于调用位置

在主线程中等待一个已经提交的异步流程,和在线程池工作线程中等待同一个线程池的新任务,不是同一回事。

尤其要警惕下面这种结构:

CompletableFuture.supplyAsync(() -> {
    // 当前任务已经占用线程池中的一个工作线程
    return CompletableFuture.supplyAsync(() -> doWork(), executor).join();
}, executor);

当线程池容量较小时,外层任务可能占满全部工作线程;内层任务虽然已经提交,却没有空闲线程执行,外层又在等待内层完成,于是形成线程饥饿。它不一定是传统意义上的锁死,但表现几乎一样:任务全部卡住。

一个可运行的并行查询示例

下面的示例模拟订单详情查询:价格和库存可以并行获取,价格结果会写入一个本地缓存。代码适用于 Java 8 及以上版本,没有依赖特定框架。

import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.ReentrantLock;

public class CompletableFutureDemo {

    static class OrderView {
        private final String sku;
        private final int price;
        private final int stock;

        OrderView(String sku, int price, int stock) {
            this.sku = sku;
            this.price = price;
            this.stock = stock;
        }

        @Override
        public String toString() {
            return "OrderView{sku='" + sku + "', price=" + price
                    + ", stock=" + stock + "}";
        }
    }

    static class ProductService implements AutoCloseable {
        private final ExecutorService ioPool = Executors.newFixedThreadPool(8);
        private final ExecutorService cpuPool = Executors.newFixedThreadPool(4);
        private final Map<String, Integer> priceCache =
                new java.util.HashMap<>();
        private final ReentrantLock cacheLock = new ReentrantLock();

        CompletableFuture<Integer> loadPrice(String sku) {
            Integer cached = readCache(sku);
            if (cached != null) {
                return CompletableFuture.completedFuture(cached);
            }

            return CompletableFuture
                    .supplyAsync(() -> queryPrice(sku), ioPool)
                    .thenApply(price -> {
                        writeCache(sku, price);
                        return price;
                    });
        }

        CompletableFuture<Integer> loadStock(String sku) {
            return CompletableFuture.supplyAsync(() -> queryStock(sku), ioPool);
        }

        CompletableFuture<OrderView> loadOrderView(String sku) {
            CompletableFuture<Integer> priceFuture = loadPrice(sku);
            CompletableFuture<Integer> stockFuture = loadStock(sku);

            return priceFuture.thenCombineAsync(
                    stockFuture,
                    (price, stock) -> new OrderView(sku, price, stock),
                    cpuPool
            );
        }

        private Integer readCache(String sku) {
            cacheLock.lock();
            try {
                return priceCache.get(sku);
            } finally {
                cacheLock.unlock();
            }
        }

        private void writeCache(String sku, int price) {
            cacheLock.lock();
            try {
                priceCache.put(sku, price);
            } finally {
                cacheLock.unlock();
            }
        }

        private int queryPrice(String sku) {
            sleep(100);
            return 199;
        }

        private int queryStock(String sku) {
            sleep(80);
            return 12;
        }

        private static void sleep(long millis) {
            try {
                Thread.sleep(millis);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                throw new IllegalStateException("task interrupted", e);
            }
        }

        @Override
        public void close() {
            ioPool.shutdown();
            cpuPool.shutdown();
            try {
                if (!ioPool.awaitTermination(3, TimeUnit.SECONDS)) {
                    ioPool.shutdownNow();
                }
                if (!cpuPool.awaitTermination(3, TimeUnit.SECONDS)) {
                    cpuPool.shutdownNow();
                }
            } catch (InterruptedException e) {
                ioPool.shutdownNow();
                cpuPool.shutdownNow();
                Thread.currentThread().interrupt();
            }
        }
    }

    public static void main(String[] args) {
        try (ProductService service = new ProductService()) {
            CompletableFuture<OrderView> future =
                    service.loadOrderView("BOOK-001");

            future.whenComplete((result, error) -> {
                if (error != null) {
                    error.printStackTrace();
                } else {
                    System.out.println(result);
                }
            }).join();
        }
    }
}

这个例子中有几个值得注意的设计点。

loadPriceloadStock 都明确使用 ioPool,说明它们属于外部等待型任务。两个 Future 通过 thenCombineAsync 合并,并且把组装工作放到 cpuPool。这里没有在线程池任务内部再创建一个 Future 后立即 join(),而是让 CompletableFuture 自己表达依赖关系。

缓存读写只在访问 HashMap 的短临界区内完成。查询价格发生在释放锁之后,因此慢速下游不会阻塞其他线程读写缓存。

为什么不建议在锁内等待异步结果

下面是一种容易被忽略的写法:

lock.lock();
try {
    return queryAsync()
            .thenApply(value -> {
                updateSharedState(value); // 这里也需要同一把锁
                return value;
            })
            .join();
} finally {
    lock.unlock();
}

调用线程持有锁并等待 queryAsync() 完成;异步回调执行到 updateSharedState 时又要获取这把锁。如果回调无法使用其他线程完成,或者后续代码还依赖当前调用线程释放资源,就可能形成锁等待。

更稳妥的方式是把“等待外部结果”和“更新共享状态”分开:

CompletableFuture<Result> future = queryAsync();

return future.thenApply(result -> {
    lock.lock();
    try {
        updateSharedState(result);
        return result;
    } finally {
        lock.unlock();
    }
});

此时异步查询期间不持有锁,只有真正修改共享状态的几行代码进入临界区。如果更新逻辑本身较复杂,还可以先在锁外计算出不可变结果,再在锁内完成一次简短替换。

常见坑与排查方式

把所有任务都放进同一个线程池

网络请求、定时任务、批处理和 CPU 计算混在一起时,一个下游接口变慢,就可能拖住整个应用。线程池至少要按任务类型隔离,并设置有界队列、合理的拒绝策略和监控指标。

用并行数量掩盖下游容量问题

allOf 可以同时发起很多任务,但它不会替你限制并发量。一次请求如果对几十个商品都发起查询,多个请求叠加后很容易打满连接池或下游限流。必要时应增加信号量、批量接口或分段提交,而不是无限扩大线程池。

异常没有被消费

CompletableFuture 的异常会沿着链路传递。如果只调用 join() 而不记录根因,日志里经常只看到 CompletionException。可以在边界位置使用 handlewhenComplete 记录业务上下文,并在真正需要返回失败时再转换异常。

future.handle((value, error) -> {
    if (error != null) {
        Throwable cause = error.getCause() == null
                ? error : error.getCause();
        System.err.println("load order failed: " + cause.getMessage());
        return fallbackValue();
    }
    return value;
});

忽略超时和取消

线程池只能管理本地任务,不能保证下游一定及时返回。Java 9 及以上可以使用 orTimeoutcompleteOnTimeout;Java 8 则需要通过定时任务自行包装超时逻辑。无论采用哪种方式,都应同时考虑连接超时、读取超时和业务总超时,避免只在 Future 层面等待。

忽略线程池关闭

业务代码中创建的线程池必须有明确的生命周期。应用关闭时调用 shutdown,必要时等待一段时间后再 shutdownNow,并正确恢复中断标记。否则测试进程可能无法退出,服务重启时也会留下未完成任务。

实践建议

  1. 为每类阻塞任务显式指定 Executor,不要把默认公共线程池当成全局业务线程池。
  2. 使用 thenCompose 表达串行依赖,使用 thenCombine 表达两个独立结果的合并,使用 allOf 表达一组任务的汇聚。
  3. 避免在线程池工作线程中等待同一个线程池提交的嵌套任务。
  4. 锁只保护共享状态的读写,不要把网络调用、数据库调用和 Future 等待放进临界区。
  5. 为线程池配置队列长度、活跃线程数、任务等待时间、拒绝次数和下游耗时监控。
  6. 对异常、超时、取消和降级路径进行测试,而不是只验证成功返回。

总结

CompletableFuture 的价值不在于把每一行代码都变成异步,而在于让任务之间的依赖关系清晰可见。线程池决定任务在哪里执行,锁决定共享状态如何安全修改,Future 组合则决定流程如何继续。三者结合时,最重要的原则是:不要在同一个执行资源上制造互相等待,也不要让锁跨越外部操作的边界。

当代码能够明确回答“这个任务在哪个线程池执行”“这里是否会阻塞”“这把锁保护哪几行数据”时,并发程序才真正具备可推理、可监控和可维护的基础。

最后修改:2026 年 08 月 30 日
如果觉得我的文章对你有用,请随意赞赏