Java

Java JUC并发编程:从概念到实战的全面指南

前言:为什么你需要系统了解JUC

如果你在Java开发中写过任何涉及并发的代码,java.util.concurrent(简称JUC)是你绕不开的核心包。它由并发编程大师Doug Lea在JDK 1.5时贡献给Java,此后历经多次迭代优化,已经成为Java并发编程的基础设施。

很多开发者对并发的理解停留在synchronized加锁、Thread启动线程的阶段。这种认知在应对高并发、高性能场景时会力不从心。JUC提供了一套经过工业级验证的并发工具,掌握它能让你写出更高效、更安全、更优雅的并发代码。

本文系统梳理JUC的核心模块,配合实战代码,帮你建立完整的JUC知识体系。

JUC解决了什么问题

synchronized的局限性

JDK 1.5之前,Java的并发控制主要依赖synchronized关键字和wait()/notify()/notifyAll()方法。这套机制能解决基本的线程安全问题,但存在明显不足:

  • 功能单一:只能实现互斥锁,无法支持读写分离、可中断等待、超时获取锁等高级场景
  • 缺乏灵活性:锁的获取和释放是隐式的,无法像显式锁那样进行细粒度控制
  • 性能瓶颈:早期synchronized的重量级实现导致上下文切换开销大(JDK 6之后通过锁升级优化大幅改善)
  • 缺少高级工具:没有信号量、倒计数门、线程池等高级并发原语

JUC的回答

JUC从底层到高层提供了一整套解决方案:

  • 基于CAS(Compare-And-Swap)实现的无锁算法,在低竞争场景下性能远超synchronized
  • 显式锁接口Lock及其实现,提供超时获取、可中断等待、公平锁等能力
  • 丰富的并发容器,在保证线程安全的同时最大化并发性能
  • 一系列同步工具类,优雅地解决多线程协调问题
  • 线程池框架,将线程管理从手工操作提升为框架化管理

核心模块详解

一、锁机制

ReentrantLock——最常用的显式锁

ReentrantLock lock = new ReentrantLock();

try {
    lock.lock();
    // 临界区代码
    System.out.println("执行业务逻辑,当前线程: " + Thread.currentThread().getName());
} finally {
    lock.unlock(); // 必须在finally中释放锁
}

相比synchronizedReentrantLock的优势:

  • tryLock():尝试获取锁,获取失败立即返回,避免死等
  • lockInterruptibly():等待锁的过程中响应中断
  • tryLock(long timeout, TimeUnit unit):超时自动释放
  • 支持公平锁:new ReentrantLock(true)

ReentrantReadWriteLock——读写分离

适用于读多写少的场景。读操作共享锁,写操作独占锁,大幅提升读密集型场景的吞吐量:

ReentrantReadWriteLock rwLock = new ReentrantReadWriteLock();

// 读操作 - 多个线程可同时持有读锁
rwLock.readLock().lock();
try {
    return data;
} finally {
    rwLock.readLock().unlock();
}

// 写操作 - 独占
rwLock.writeLock().lock();
try {
    data = newValue;
} finally {
    rwLock.writeLock().unlock();
}

StampedLock——JDK 8的优化

在读写锁基础上增加了乐观读模式,读操作无需真正加锁,性能更高:

StampedLock sl = new StampedLock();
long stamp = sl.tryOptimisticRead(); // 不加锁,返回stamp
// 读取数据
int value = data;
if (!sl.validate(stamp)) {           // 检查期间是否有写操作
    stamp = sl.readLock();           // 退化为悲观读锁
    try { value = data; } finally { sl.unlockRead(stamp); }
}

注意:StampedLock不可重入,不支持条件变量,使用时需谨慎。

二、原子类——CAS的封装

原子类底层基于CAS(Compare-And-Swap)指令实现,无需加锁即可完成线程安全操作:

AtomicInteger count = new AtomicInteger(0);
count.incrementAndGet();       // 原子自增
count.compareAndSet(0, 1);    // CAS操作

AtomicReference<User> ref = new AtomicReference<>(new User("Tom"));
ref.compareAndSet(oldUser, newUser);

LongAdder vs AtomicLong:在高竞争场景下,AtomicLong的CAS自旋会导致大量线程空转。LongAdder通过分段累加、最终汇总的策略大幅降低竞争:

LongAdder adder = new LongAdder();
adder.increment();    // 分散到多个Cell中累加
long total = adder.sum(); // 需要时汇总

统计计数器场景优先选LongAdder,需要精确原子值时用AtomicLong

三、并发容器

ConcurrentHashMap

JDK 8之后,ConcurrentHashMap放弃了分段锁设计,采用CAS + synchronized锁定桶头节点的方式,粒度更细、并发度更高:

ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<>();
map.put("key", 1);
// 原子操作:key不存在时put
map.putIfAbsent("key", 2);
// 原子操作:key存在时更新
map.computeIfPresent("key", (k, v) -> v + 1);

重要提示:ConcurrentHashMap.size()返回的是近似值,因为在获取大小的瞬间可能有其他线程在修改。需要精确计数时使用mappingCount()(返回long,减少溢出风险)或自行加锁统计。

CopyOnWriteArrayList

写时复制策略:每次写操作都复制整个底层数组,读操作无锁。适用于读极多写极少的场景(如配置列表、监听器列表):

CopyOnWriteArrayList<String> list = new CopyOnWriteArrayList<>();
list.add("item1");
// 读操作不加锁,写操作自动复制

注意:写操作开销大,且数据一致性是最终一致的,不适合写频繁的场景。

BlockingQueue家族

阻塞队列是生产者-消费者模式的核心组件:

队列类型特点
ArrayBlockingQueue有界数组队列,必须指定容量
LinkedBlockingQueue可选有界链表队列,默认Integer.MAX_VALUE
PriorityBlockingQueue优先级排序的无界队列
SynchronousQueue零容量队列,生产者直接交给消费者
DelayQueue延迟队列,元素到期后才能取出

四、同步工具类

CountDownLatch——倒计时门闩

一个线程等待其他线程完成后再执行:

CountDownLatch latch = new CountDownLatch(3);

// 三个工作线程
for (int i = 0; i < 3; i++) {
    executor.submit(() -> {
        doWork();
        latch.countDown(); // 完成一个任务,计数减1
    });
}

// 主线程等待所有任务完成
latch.await();
System.out.println("所有任务完成,继续执行");

CyclicBarrier——循环栅栏

一组线程互相等待,全部到达屏障点后一起执行下一步,且可重复使用:

CyclicBarrier barrier = new CyclicBarrier(3, () -> {
    System.out.println("所有线程到达屏障,执行汇总操作");
});

// 每个线程到达屏障点后等待
barrier.await(); // 阻塞直到3个线程都到达

CountDownLatch的核心区别:CountDownLatch是一次性的,CyclicBarrier可重置重用。

Semaphore——信号量/限流器

控制同时访问资源的线程数量,经典限流场景:

Semaphore semaphore = new Semaphore(5); // 允许5个线程同时访问

semaphore.acquire(); // 获取许可,无可用许可时阻塞
try {
    accessResource();
} finally {
    semaphore.release(); // 释放许可
}

Phaser——阶段器

JDK 7引入,是CountDownLatchCyclicBarrier的增强版,支持动态注册参与者、多阶段同步:

Phaser phaser = new Phaser(1); // 初始注册1个参与者
phaser.register();  // 动态注册新参与者
phaser.arriveAndAwaitAdvance(); // 到达并等待其他参与者
phaser.arriveAndDeregister();   // 到达并注销

五、线程池——Executor框架

ThreadPoolExecutor

生产环境中永远不要使用Executors工厂方法创建线程池(Executors.newFixedThreadPoolExecutors.newCachedThreadPool的队列或线程数无界,存在OOM风险),应该显式构造:

ThreadPoolExecutor executor = new ThreadPoolExecutor(
    4,                          // 核心线程数
    8,                          // 最大线程数
    60L, TimeUnit.SECONDS,      // 非核心线程空闲存活时间
    new LinkedBlockingQueue<>(100), // 有界任务队列
    new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略
);

拒绝策略选择

策略行为
AbortPolicy抛出RejectedExecutionException(默认)
CallerRunsPolicy调用者线程执行任务(降级方案)
DiscardPolicy静默丢弃
DiscardOldestPolicy丢弃队列最老任务

参数配置经验

  • CPU密集型任务:核心线程数 ≈ CPU核心数 + 1
  • IO密集型任务:核心线程数 ≈ CPU核心数 × 2 ~ 4
  • 混合场景:根据任务IO耗时占比调整

ForkJoinPool

适合可以递归拆分的计算密集型任务,采用work-stealing策略:

ForkJoinPool pool = new ForkJoinPool();
long result = pool.invoke(new RecursiveTask<Long>() {
    @Override
    protected Long compute() {
        if (n <= THRESHOLD) {
            return directCompute(n);
        }
        RecursiveTask<Long> left = new SubTask(begin, mid);
        RecursiveTask<Long> right = new SubTask(mid + 1, end);
        left.fork(); // 异步执行
        long rightResult = right.compute(); // 当前线程执行
        long leftResult = left.join();       // 获取结果
        return leftResult + rightResult;
    }
});

六、CompletableFuture——异步编程利器

CompletableFuture是JDK 8引入的异步编程框架,支持任务编排、组合、回调,比Future强大得多:

// 异步执行
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    return queryFromDB();
});

// 链式处理
future.thenApply(result -> transform(result))
      .thenAccept(finalResult -> save(finalResult))
      .exceptionally(ex -> { log(ex); return null; });

// 组合多个异步任务
CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> queryUser());
CompletableFuture<List<Order>> f2 = CompletableFuture.supplyAsync(() -> queryOrders());

CompletableFuture<Result> combined = f1.thenCombine(f2, (user, orders) -> {
    return new Result(user, orders);
});

典型应用场景实战

生产者-消费者模型

BlockingQueue<Task> queue = new ArrayBlockingQueue<>(100);

// 生产者
executor.submit(() -> {
    while (!stopped) {
        Task task = generateTask();
        queue.put(task); // 队列满时阻塞
    }
});

// 消费者
executor.submit(() -> {
    while (!stopped) {
        Task task = queue.take(); // 队列空时阻塞
        process(task);
    }
});

限流场景

public class RateLimiter {
    private final Semaphore semaphore;

    public RateLimiter(int maxConcurrency) {
        this.semaphore = new Semaphore(maxConcurrency);
    }

    public <T> T execute(Supplier<T> action) throws InterruptedException {
        semaphore.acquire();
        try {
            return action.get();
        } finally {
            semaphore.release();
        }
    }
}

异步编排:并行查询多个服务

CompletableFuture<User> userFuture = CompletableFuture.supplyAsync(userService::getUser);
CompletableFuture<List<Order>> orderFuture = CompletableFuture.supplyAsync(orderService::getOrders);
CompletableFuture<Account> accountFuture = CompletableFuture.supplyAsync(accountService::getAccount);

// 等待全部完成,组合结果
CompletableFuture<Dashboard> dashboard = CompletableFuture
    .allOf(userFuture, orderFuture, accountFuture)
    .thenApply(v -> new Dashboard(
        userFuture.join(),
        orderFuture.join(),
        accountFuture.join()
    ));

常见陷阱与注意事项

死锁风险与锁顺序

多个锁的获取顺序不一致会导致死锁。解决办法:始终按固定顺序获取锁,或使用tryLock超时机制兜底。

线程池的坑

  • Executors.newFixedThreadPool()使用无界LinkedBlockingQueue,任务堆积会导致OOM
  • Executors.newCachedThreadPool()最大线程数为Integer.MAX_VALUE,可能创建大量线程
  • 任务提交时线程池已关闭会抛RejectedExecutionException
  • shutdown()后不再接受新任务,但已提交的任务会继续执行

ConcurrentHashMap的size()

size()在高并发下是非精确的。如果需要精确计数,使用mappingCount()或者在统计时加外部锁。

volatile的局限

volatile保证可见性和有序性,但不保证原子性。i++这样的复合操作即使用了volatile也不安全,应改用AtomicInteger

CompletableFuture异常处理

异步链中任何一个环节抛异常,如果不处理,异常会被静默吞掉。务必在链尾加exceptionally()handle()

future.thenApply(this::process)
      .exceptionally(ex -> {
          log.error("处理失败", ex);
          return defaultValue;
      });

ThreadLocal内存泄漏

ThreadLocal在使用线程池时容易泄漏——线程被复用但ThreadLocal值未清理。必须在finally块中调用remove()

ThreadLocal<UserContext> context = new ThreadLocal<>();
try {
    context.set(userContext);
    doWork();
} finally {
    context.remove(); // 必须清理
}

JUC vs synchronized:怎么选

维度synchronizedJUC (ReentrantLock等)
使用复杂度简单,自动管理锁需手动加锁/释放
灵活性高(超时、中断、公平锁)
性能(低竞争)JDK 6后优化,差距不大CAS路径更快
性能(高竞争)线程挂起/唤醒开销大支持自旋,可调优
读写分离不支持ReentrantReadWriteLock
条件变量单条件(wait/notify)多条件(Condition)

建议:简单同步场景用synchronized(代码简洁、不易出错);需要高级特性(超时、可中断、读写锁、多条件)时用JUC。

Java 21+ 虚拟线程对JUC的影响

Java 21正式引入虚拟线程(Virtual Threads),这是Java并发模型的一次重大变革。虚拟线程是轻量级线程,由JVM调度而非操作系统,创建和切换成本极低。

虚拟线程改变了什么

  • 线程池不再是必须:虚拟线程创建成本几乎为零,Executors.newVirtualThreadPerTaskExecutor()可以替代大部分ThreadPoolExecutor场景
  • IO密集型应用收益巨大:不再受平台线程数限制,可以一个请求一个线程
  • ForkJoinPool地位下降:虚拟线程本身就是work-stealing的

虚拟线程没有改变什么

  • synchronizedReentrantLock等锁机制在虚拟线程下依然有效
  • ConcurrentHashMap等并发容器依然是线程安全的首选
  • CountDownLatchSemaphore等同步工具照常使用

需要注意的点

  • 虚拟线程使用ForkJoinPool调度,synchronized块内的pin操作会导致载体线程被占用,建议改用ReentrantLock
  • ThreadLocal在虚拟线程场景下泄漏风险更大(线程数远超平台线程),应优先使用ScopedValue(预览特性)
  • JUC的大部分工具在虚拟线程环境下行为不变,但ForkJoinPool的使用模式需要调整

虚拟线程不会取代JUC,而是让JUC的部分组件(主要是线程池相关)变得可选。锁机制、并发容器、同步工具这些基础设施依然重要。

结语

JUC是Java并发编程的核心武器库。从锁机制到原子类,从并发容器到同步工具,从线程池到异步编排,它提供了一套完整的解决方案。

学习JUC不是死记API,而是理解每个工具背后的设计思想和适用场景。ReentrantLock解决的是锁的灵活控制问题,ConcurrentHashMap解决的是高并发下的数据一致性,CompletableFuture解决的是异步任务的优雅编排。

最后,虚拟线程的到来是好消息而非威胁——它让你在高并发场景下有更多选择,而JUC作为经过工业验证的基础设施,依然是你工具箱中不可或缺的一部分。

更早的文章

SSE流式响应:前端取消请求后,服务端还在烧Token吗?

欢迎在评论区留下您的见解~