10. Java 高并发与并发工具包

深入 Java 并发编程,讲解 JUC 线程池、CAS 乐观锁、AQS 框架、锁升级与并发集合实战方案

Java 的高并发能力源于 JUC(java.util.concurrent)包和 JVM 底层的内存模型与锁机制。理解 CAS、AQS、线程池与锁升级的全链路,是设计高吞吐系统的关键。

1. 并发基础:JMM 与 happens-before

1.1 重排序与可见性

JMM 允许编译器和 CPU 对指令重排序(只要不改变单线程语义)。单例模式的双重检查锁定必须在实例字段上声明 volatile

public class Singleton {
    private volatile static Singleton instance;
    
    public static Singleton getInstance() {
        if (instance == null) {                // 第一次检查(无锁)
            synchronized (Singleton.class) {
                if (instance == null) {        // 第二次检查(有锁)
                    instance = new Singleton(); // 非原子操作:分配内存→初始化→引用赋值
                }
            }
        }
        return instance;
    }
}

volatile 的作用:

  • 可见性:写操作立即刷入主内存,读操作从主内存刷新
  • 禁止指令重排序new Singleton() 的三步不会重排,防止返回未完全初始化的对象

1.2 happens-before 八大规则

规则说明
程序次序规则同一线程中前面的操作 happens-before 后面的操作
管程锁定规则unlock happens-before 后续对同一锁的 lock
volatile 规则写 happens-before 后续读
线程启动规则start() happens-before 线程中所有操作
线程终止规则线程中所有操作 happens-before join() 返回
中断规则interrupt() happens-before 检测中断状态
对象终结规则构造函数执行 happens-before finalize()
传递性A→B 且 B→C ⟹ A→C

2. CAS 与原子操作

2.1 CAS 原理

// Unsafe 底层实现(伪代码)
public final native boolean compareAndSwapInt(
    Object o, long offset, int expected, int x);

CAS 存在三大问题:

  1. ABA 问题:值由 A→B→A,CAS 误认为未变化 → 使用 AtomicStampedReference(版本号)
  2. 自旋开销:长时间失败导致 CPU 空转 → 自适应自旋或退化为锁
  3. 只能保证单个变量原子性:多变量需用 AtomicReference 封装或锁

2.2 LongAdder 替代 AtomicLong

// 高并发下 AtomicLong 的 CAS 竞争激烈
private final AtomicLong count = new AtomicLong(0);
count.incrementAndGet();

// LongAdder:分散热点,低争用时性能接近 AtomicLong,高争用时提升 10 倍
private final LongAdder adder = new LongAdder();
adder.increment();
long total = adder.sum(); // 非精确,适合统计计数

LongAdder 内部维护 Cell[] 数组,线程映射到不同 Cell,最后求和。

3. AQS 框架解析

3.1 AQS 核心结构

public abstract class AbstractQueuedSynchronizer
    extends AbstractOwnableSynchronizer {
    
    volatile int state;           // 同步状态
    volatile Node head;           // CLH 队列头
    volatile Node tail;           // CLH 队列尾
    
    // 子类实现:独占式
    protected boolean tryAcquire(int arg);
    protected boolean tryRelease(int arg);
    
    // 子类实现:共享式
    protected int tryAcquireShared(int arg);
    protected boolean tryReleaseShared(int arg);
}

AQS 队列是 CLH 变体的双向 FIFO 队列

  • 每个节点持有前驱引用,释放时唤醒后继
  • 通过 LockSupport.park/unpark 阻塞/唤醒线程,避免自旋浪费 CPU

3.2 ReentrantLock 实现剖析

public class ReentrantLock implements Lock {
    private final Sync sync;
    
    abstract static class Sync extends AQS {
        // 公平锁:检查队列是否有前驱等待节点
        final boolean nonfairTryAcquire(int acquires) {
            final Thread current = Thread.currentThread();
            int c = getState();
            if (c == 0) {
                if (compareAndSetState(0, acquires)) { // 直接 CAS 抢锁
                    setExclusiveOwnerThread(current);
                    return true;
                }
            } else if (current == getExclusiveOwnerThread()) {
                setState(c + acquires); // 重入
                return true;
            }
            return false;
        }
    }
    
    static final class FairSync extends Sync {
        protected final boolean tryAcquire(int acquires) {
            final Thread current = Thread.currentThread();
            int c = getState();
            if (c == 0) {
                if (!hasQueuedPredecessors() &&  // 公平性关键:检查是否有排队者
                    compareAndSetState(0, acquires)) {
                    setExclusiveOwnerThread(current);
                    return true;
                }
            }
            // ... 重入逻辑
        }
    }
}
特性ReentrantLocksynchronized
实现APIJVM 监视器锁
可重入
公平锁✅ 可选
可中断lockInterruptibly
超时等待tryLock(timeout)
条件变量✅ 多 Conditionwait/notify
性能JDK6+ 基本持平同左

4. 锁升级与 Mark Word

4.1 对象头结构(64 位 JVM)

|--------------------------------------------------|
| Mark Word (64 bits)                               |
|--------------------------------------------------|
| 锁状态    | 56 bits (分代年龄/hashCode/锁记录等) |
|--------------------------------------------------|
| Class Pointer (64 bits, 压缩后 32 bits)          |
|--------------------------------------------------|
| Array Length (仅数组对象)                          |
|--------------------------------------------------|

4.2 锁升级路径

无锁 → 偏向锁 → 轻量级锁 → 重量级锁
阶段条件机制
偏向锁单线程反复进入同步块Mark Word 存储线程 ID,CAS 替换,失败则撤销
轻量级锁多线程交替执行(无竞争)栈帧创建 Lock Record,CAS 替换 Mark Word,自旋 10 次
重量级锁竞争剧烈或自旋失败膨胀为 ObjectMonitor,线程进入 cxq/EntryList 阻塞

偏向锁在 JDK 15 默认关闭(-XX:-UseBiasedLocking),多线程环境下偏向→撤销的开销可能大于收益。

4.3 自旋锁适应化

-XX:+UseSpinning          # JDK6+ 默认开启
-XX:PreBlockSpin=10       # 默认自旋 10 次

自适应自旋:根据上一个线程在临界区的执行时长,动态调整当前线程的自旋次数。

5. 线程池深度解析

5.1 ThreadPoolExecutor 参数

new ThreadPoolExecutor(
    4,                      // corePoolSize: 常驻核心线程数
    8,                      // maximumPoolSize: 最大线程数
    60L, TimeUnit.SECONDS,  // keepAliveTime: 非核心线程空闲存活时间
    new LinkedBlockingQueue<>(100), // workQueue: 任务队列
    new ThreadFactory() {...},      // threadFactory: 自定义线程名、守护状态
    new ThreadPoolExecutor.CallerRunsPolicy() // rejectionHandler
);

5.2 任务提交流程

1. 当前线程数 < corePoolSize → 创建新核心线程处理
2. 当前线程数 ≥ corePoolSize → 任务进入 workQueue
3. workQueue 满 → 创建非核心线程(不超过 maximumPoolSize)
4. 线程数达到 maximumPoolSize 且队列满 → 触发 RejectedExecutionHandler

5.3 拒绝策略

策略行为场景
AbortPolicy抛出异常默认,快速失败
CallerRunsPolicy由提交线程(主线程)执行降低提交速度,自我保护
DiscardPolicy静默丢弃允许丢弃非关键任务
DiscardOldestPolicy丢弃最旧任务,重试提交追求最新数据

5.4 线程池规范实践

public class ThreadPoolConfig {
    
    @Bean("ioThreadPool")
    public ThreadPoolExecutor ioThreadPool() {
        return new ThreadPoolExecutor(
            2, 4, 60, TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(1000),
            new ThreadFactoryBuilder().setNameFormat("io-pool-%d").build(),
            new ThreadPoolExecutor.CallerRunsPolicy()
        );
    }
    
    @Bean("cpuThreadPool")
    public ThreadPoolExecutor cpuThreadPool() {
        int cpu = Runtime.getRuntime().availableProcessors();
        return new ThreadPoolExecutor(
            cpu, cpu * 2, 60, TimeUnit.SECONDS,
            new ArrayBlockingQueue<>(100),
            new ThreadFactoryBuilder().setNameFormat("cpu-pool-%d").build(),
            new ThreadPoolExecutor.AbortPolicy()
        );
    }
}

// 优雅关闭
public void gracefulShutdown(ThreadPoolExecutor pool) {
    pool.shutdown(); // 停止接收新任务
    try {
        if (!pool.awaitTermination(60, TimeUnit.SECONDS)) {
            pool.shutdownNow(); // 强制中断
        }
    } catch (InterruptedException e) {
        pool.shutdownNow();
    }
}

6. 并发集合

6.1 ConcurrentHashMap 演进

JDK 7:Segment 分段锁(默认 16 段)→ 锁粒度为段,高并发下段内仍有竞争。

JDK 8:CAS + synchronized(Node 级锁)+ 红黑树

// put 流程(JDK 8)
final V putVal(K key, V value, boolean onlyIfAbsent) {
    if (key == null || value == null) throw new NullPointerException();
    int hash = spread(key.hashCode());
    int binCount = 0;
    for (Node<K,V>[] tab = table;;) {
        Node<K,V> f; int n, i, fh;
        if (tab == null || (n = tab.length) == 0)
            tab = initTable();           // 懒初始化
        else if ((f = tabAt(tab, i = (n - 1) & hash)) == null) {
            if (casTabAt(tab, i, null, new Node<K,V>(hash, key, value)))
                break;                    // 无竞争:CAS 直接插入
        }
        else if ((fh = f.hash) == MOVED)
            tab = helpTransfer(tab, f);  // 协助扩容
        else {
            synchronized (f) {           // 有竞争:锁头节点
                // 链表遍历/红黑树插入
            }
        }
    }
    addCount(1L, binCount);              // LongAdder 统计 size
}

6.2 CopyOnWriteArrayList

// 读无锁,写时复制整个数组。适合读多写极少
private transient volatile Object[] array;

public boolean add(E e) {
    final ReentrantLock lock = this.lock;
    lock.lock();
    try {
        Object[] elements = getArray();
        int len = elements.length;
        Object[] newElements = Arrays.copyOf(elements, len + 1);
        newElements[len] = e;
        setArray(newElements); // volatile 写,保证可见性
        return true;
    } finally {
        lock.unlock();
    }
}

注意:CopyOnWrite 仅保证最终一致性,迭代器可能读到旧数据。适合事件监听列表等场景。

7. JUC 工具类

7.1 CountDownLatch

CountDownLatch latch = new CountDownLatch(3);

// 主线程等待
latch.await();

// 子线程计数减一
latch.countDown();

应用:多数据源并行加载后聚合结果,或等待所有服务启动完成。

7.2 CyclicBarrier

CyclicBarrier barrier = new CyclicBarrier(3, () -> {
    System.out.println("所有线程就绪,开始下一阶段");
});

barrier.await(); // 计数达到 3 时触发 barrierAction,然后重置(可循环使用)

应用:分阶段并行计算,如 MapReduce 的 Map 阶段完成后统一进入 Reduce。

7.3 Semaphore

Semaphore semaphore = new Semaphore(10); // 限流 10 并发

if (semaphore.tryAcquire(100, TimeUnit.MILLISECONDS)) {
    try {
        // 执行业务
    } finally {
        semaphore.release();
    }
}

应用:数据库连接池、API 接口限流。

7.4 CompletableFuture 并发编排

CompletableFuture<Integer> futureA = CompletableFuture
    .supplyAsync(() -> fetchOrderPrice(orderId), orderPool);
CompletableFuture<Integer> futureB = CompletableFuture
    .supplyAsync(() -> fetchUserCoupon(userId), couponPool);

CompletableFuture<Integer> total = futureA
    .thenCombine(futureB, (price, coupon) -> price - coupon)
    .orTimeout(3, TimeUnit.SECONDS)
    .exceptionally(ex -> {
        log.error("计算失败", ex);
        return price; // 降级:不使用优惠券
    });

8. Fork/Join 框架

public class ArraySumTask extends RecursiveTask<Long> {
    private static final int THRESHOLD = 10000;
    private final int[] array;
    private final int start, end;
    
    @Override
    protected Long compute() {
        if (end - start <= THRESHOLD) {
            long sum = 0;
            for (int i = start; i < end; i++) sum += array[i];
            return sum;
        }
        int mid = (start + end) / 2;
        ArraySumTask left = new ArraySumTask(array, start, mid);
        ArraySumTask right = new ArraySumTask(array, mid, end);
        left.fork();  // 异步执行左半
        long rightResult = right.compute(); // 当前线程执行右半
        long leftResult = left.join(); // 等待左半结果
        return leftResult + rightResult;
    }
}

Fork/Join 使用 工作窃取(Work Stealing):每个线程维护双端队列,空闲线程从其他线程队列尾部窃取任务。

9. 常见并发陷阱

陷阱现象解决
线程池共享变量未 volatile读不到最新值volatile 或原子类
HashMap 并发 put死循环(JDK7)或数据丢失改用 ConcurrentHashMap
线程池 rejected 无监控任务静默丢失自定义拒绝策略记录日志
SimpleDateFormat 共享格式化结果错乱使用 DateTimeFormatter(线程安全)或 ThreadLocal
CompletableFuture 默认线程池ForkJoinPool.common() 可能被阻塞始终指定自定义线程池

总结

Java 高并发编程的核心是理解 JMM happens-before 规则CAS 乐观锁的局限AQS 队列同步框架。实际工程中:

  1. 优先使用 java.util.concurrent 工具类而非 wait/notify
  2. 线程池必须自定义参数,禁止 Executors.newFixedThreadPool()(无界队列导致 OOM)
  3. 锁按固定顺序获取,避免死锁
  4. 线上必须暴露线程池指标(活跃线程、队列堆积、拒绝次数)
  5. JDK 21+ 引入轻量级线程,但 JUC 知识体系仍是并发编程的基石

掌握这些内容后,可以从容应对高并发系统的线程安全设计与性能调优。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java-enterprise」更多文章

  1. 限流算法深度解析:令牌桶、漏桶与滑动窗口计数
  2. Java 代码质量:SonarQube、Checkstyle 与 SpotBugs 工程化实践
  3. Spring IoC 容器与依赖注入原理深度剖析