Java并发工具类案例

wen java案例 1

Java并发工具类实战:从CountDownLatch到CompletableFuture的五大核心场景解析

目录导读

  1. 为什么你需要掌握并发工具类?
  2. CountDownLatch — 多任务并行等待
  3. CyclicBarrier — 循环栅栏的协同突破
  4. Semaphore — 限流与资源池控制
  5. ConcurrentHashMap 与 ConcurrentSkipListMap — 高并发Map选型
  6. CompletableFuture — 异步编排的“银弹”
  7. 高频面试问答(Q&A)
  8. 如何选择正确的并发工具

为什么你需要掌握并发工具类?

在Java多线程开发中,synchronizedLock解决了互斥问题,但面对任务协同资源限制异步流水线等复杂场景,直接使用底层API会陷入繁琐且易错的等待/通知逻辑,Java并发包(java.util.concurrent)提供了一组高抽象级别的工具类,它们封装了状态判断与线程阻塞唤醒,让开发者用最少代码实现可靠并发协作。

Java并发工具类案例

核心痛点:生产环境中的线程池耗尽、竞态条件、死锁,90%源于对工具类选择不当或误用。


场景一:CountDownLatch — 多任务并行等待

案例:某电商大促需要同时拉取用户信息、库存、优惠券三个接口,全部完成后渲染页面。

CountDownLatch latch = new CountDownLatch(3);
ExecutorService pool = Executors.newFixedThreadPool(3);
pool.submit(() -> { try { fetchUser(); } finally { latch.countDown(); } });
pool.submit(() -> { try { fetchStock(); } finally { latch.countDown(); } });
pool.submit(() -> { try { fetchCoupon(); } finally { latch.countDown(); } });
latch.await(3, TimeUnit.SECONDS); // 主线程等待,超时避免挂死
renderPage();

关键点

  • 计数器不可复用(与CyclicBarrier区别)。
  • 必须finally中调用countDown(),防止子任务异常导致主线程永久阻塞。
  • await(timeout)是生产必备,防止依赖服务无响应。

场景二:CyclicBarrier — 循环栅栏的协同突破

案例:一个大数据报表任务,分成4个批次(Partition)统计,每批次中4个线程各自计算一部分,计算完成后合并。

int batchCount = 4;
CyclicBarrier barrier = new CyclicBarrier(4, () -> System.out.println("本批次合并完成"));
for (int i = 0; i < batchCount; i++) {
    int batch = i;
    for (int j = 0; j < 4; j++) {
        int part = j;
        pool.submit(() -> {
            // 模拟计算
            compute(batch, part);
            try {
                barrier.await(); // 每批次4个线程相互等待
            } catch (InterruptedException | BrokenBarrierException e) {
                Thread.currentThread().interrupt();
            }
        });
    }
}

与CountDownLatch对比: | 维度 | CountDownLatch | CyclicBarrier | |------|---------------|---------------| | 等待方式 | 一个/多个线程等待其他线程完成 | 所有线程互相等待,同时被唤醒 | | 复用性 | 不能复用 | 可以reset()后再次使用 | | 应用场景 | 事件完成通知 | 固定数量线程反复协同 |

陷阱:若一个线程在await()之前异常退出,其余线程会抛出BrokenBarrierException,需处理中断。


场景三:Semaphore — 限流与资源池控制

案例:数据库连接池仅允许10个并发连接,超出线程必须等待。

Semaphore semaphore = new Semaphore(10, true); // 公平模式
public Connection getConnection() {
    try {
        if (semaphore.tryAcquire(500, TimeUnit.MILLISECONDS)) {
            try {
                return doGetConnection(); // 核心获取逻辑
            } finally {
                semaphore.release();
            }
        } else {
            throw new TimeoutException("获取连接超时");
        }
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
        return null;
    }
}

最佳实践

  • tryAcquire(timeout)优于acquire(),避免无限等待。
  • Semaphore支持非对称释放(一个线程获取,另一线程释放),适用于生产者-消费者模型。
  • 公平模式(new Semaphore(permits, true))可防止线程饥饿。

场景四:ConcurrentHashMap 与 ConcurrentSkipListMap — 高并发Map选型

案例:缓存活跃用户ID(随机访问) vs 维护实时排名(有序遍历)。

工具类 数据结构 并发性能 有序性 适用场景
ConcurrentHashMap 哈希表+红黑树 极高(CAS+分段锁) 普通KV缓存、计数
ConcurrentSkipListMap 跳表 高(线程安全的排序) 键有序 排行榜、区间查询

核心细节

  • ConcurrentHashMapget()完全无锁,put()仅在桶为空时CAS,桶非空时锁住单个桶(JDK8+)。
  • 禁止使用size()在并发下精确控制业务(它是个估计值)。
  • 为复合操作(如putIfAbsent后累加)使用compute/merge方法,避免get+put的竞态条件。

场景五:CompletableFuture — 异步编排的“银弹”

案例:一个订单流程:查询用户 →(并行)查余额 & 查风控 → 全部完成后创建订单 → 失败回退。

CompletableFuture<User> userFuture = CompletableFuture.supplyAsync(() -> getUser());
CompletableFuture<Double> balanceFuture = userFuture
    .thenCompose(user -> CompletableFuture.supplyAsync(() -> getBalance(user)));
CompletableFuture<Boolean> riskFuture = userFuture
    .thenCompose(user -> CompletableFuture.supplyAsync(() -> checkRisk(user)));
CompletableFuture<Void> orderFuture = balanceFuture
    .thenCombine(riskFuture, (balance, risk) -> 
        risk ? -1 : createOrder(balance))
    .exceptionally(ex -> { // 异常回退
        log.error("订单创建失败", ex);
        return -2;
    })
    .thenAccept(orderId -> System.out.println("结果:" + orderId));
orderFuture.join(); // 阻塞等待整个异步链路结束(仅测试用)

关键API理解

  • supplyAsync:异步执行有返回值任务。
  • thenCompose:扁平化异步依赖(类似flatMap),避免嵌套Future。
  • thenCombine:两个异步结果合并。
  • exceptionally:捕获上游任何异常并恢复链路。
  • 默认线程池(ForkJoinPool.commonPool)在IO密集型场景下可能成为瓶颈,务必自定义线程池作为第二参数

高频面试问答(Q&A)

Q1: CountDownLatch和CyclicBarrier都能等待线程,你能现场写一个区别点吗? A: 核心区别是应用语义,CountDownLatch是“裁判等运动员”——一个/几个裁判线程等所有运动员跑完;CyclicBarrier是“几个运动员互相等”——每批人都到齐才一起出发,且下一批可复用,另外CountDownLatch的计数只能减,CyclicBarrier可重置。

Q2: Semaphore可以用来实现互斥锁吗? A: 可以。new Semaphore(1)就等价于一个非重入的互斥锁,但要注意非重入——同一线程再次获取会死锁,所以若需要可重入且简单锁,优先用ReentrantLock

Q3: CompletableFuture中thenApplythenCompose区别? A: thenApply对结果同步映射(返回一个值),thenCompose返回一个CompletableFuture(用于扁平化解依赖链),使用thenCompose是为了避免CompletableFuture<CompletableFuture<T>>的嵌套结构。

Q4: 并发下如何安全地对ConcurrentHashMap做“值累加”? A: 使用map.compute(key, (k, v) -> v == null ? 1 : v + 1),这是一个原子复合操作,禁止用get + put组合。


如何选择正确的并发工具

  • 等待一组任务完成CountDownLatch
  • 多个线程互相等待到齐后同时突破CyclicBarrier
  • 限制并发访问数Semaphore
  • 高并发KV存储(无需排序)ConcurrentHashMap
  • 高并发且需要有序访问ConcurrentSkipListMap
  • 复杂异步流水线、依赖编排CompletableFuture

最后提示:工具类只是手段,真正的并发安全取决于你对共享状态边界的划分,推荐遵循“不可变对象优先 + 用工具类管理协作”的模式,避免将volatile与锁混用导致的逻辑混乱,在常规业务中,优先考虑将这些工具类封装为线程安全的服务组件,并配合监控线程池指标(活跃线程数、队列深度),才能确保系统在高并发下稳定运行。

抱歉,评论功能暂时关闭!