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 存在三大问题:
- ABA 问题:值由 A→B→A,CAS 误认为未变化 → 使用
AtomicStampedReference(版本号) - 自旋开销:长时间失败导致 CPU 空转 → 自适应自旋或退化为锁
- 只能保证单个变量原子性:多变量需用
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;
}
}
// ... 重入逻辑
}
}
}
| 特性 | ReentrantLock | synchronized |
|---|---|---|
| 实现 | API | JVM 监视器锁 |
| 可重入 | ✅ | ✅ |
| 公平锁 | ✅ 可选 | ❌ |
| 可中断 | ✅ lockInterruptibly | ❌ |
| 超时等待 | ✅ tryLock(timeout) | ❌ |
| 条件变量 | ✅ 多 Condition | 单 wait/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 队列同步框架。实际工程中:
- 优先使用
java.util.concurrent工具类而非wait/notify - 线程池必须自定义参数,禁止
Executors.newFixedThreadPool()(无界队列导致 OOM) - 锁按固定顺序获取,避免死锁
- 线上必须暴露线程池指标(活跃线程、队列堆积、拒绝次数)
- JDK 21+ 引入轻量级线程,但 JUC 知识体系仍是并发编程的基石
掌握这些内容后,可以从容应对高并发系统的线程安全设计与性能调优。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。