发布于2026-06-30 阅读(0)
扫一扫,手机访问
在多核CPU早已成为标配的今天,并发编程已经不是什么进阶技能,而是后端开发者必须掌握的看家本领了。这篇文章我们系统梳理一下Ja va并发编程中几个最核心的模块:锁机制、并发容器、阻塞队列以及异步编程。不光讲原理,也会配上实战代码,帮大家在高并发场景下做出更明智的技术选型,少踩一些坑。

synchronized是Ja va原生的隐式锁,基于JVM实现,加锁解锁全自动;而ReentrantLock是JDK提供的显式锁,基于AQS(AbstractQueuedSynchronizer)框架,需要手动控制。
对比维度 | synchronized | ReentrantLock |
实现层面 | JVM 层面 | JDK API 层面 |
锁释放 | 自动释放 | 必须手动 unlock () |
可中断 | 不支持 | 支持 lockInterruptibly () |
可超时 | 不支持 | 支持 tryLock (timeout) |
公平锁 | 仅非公平 | 支持公平 / 非公平 |
条件变量 | 仅 1 个等待队列 | 支持多个 Condition |
当然,synchronized在JDK6之后做了大量优化,像偏向锁、轻量级锁、自旋锁这些,性能已经大幅提升了。不过,在需要灵活控制锁行为、或者需要更精细的并发控制时,ReentrantLock依然是绕不开的选择。
可重入锁的意思是:同一个线程可以多次获取同一把锁而不会把自己给卡死。ReentrantLock通过AQS的state状态变量和exclusiveOwnerThread来实现这一点。
具体来说,线程首次获取锁时,state从0变成1,同时记录下当前持有线程是谁。同一线程再次获取锁时,state继续累加;释放时则递减,直到state归零,锁才真正释放。
// ReentrantLock.NonfairSync.tryAcquire()核心逻辑
final boolean nonfairTryAcquire(int acquires) {
final Thread current = Thread.currentThread();
int c = getState();
if (c == 0) {
// 锁空闲,CAS尝试获取
if (compareAndSetState(0, acquires)) {
setExclusiveOwnerThread(current);
return true;
}
}
else if (current == getExclusiveOwnerThread()) {
// 同一线程重入,state累加
int nextc = c + acquires;
if (nextc < 0) throw new Error("Maximum lock count exceeded");
setState(nextc);
return true;
}
return false;
}
可重入性避免了递归调用或同一线程反复获取锁导致的死锁,这是它最核心的价值所在。
synchronized获取锁时是不可中断的——线程会一直傻等,直到拿到锁为止。ReentrantLock则通过lockInterruptibly()支持中断响应,给阻塞中的线程一个“反悔”的机会。
当调用lockInterruptibly()时,如果线程在等待队列中被中断了,它会直接抛出InterruptedException,不再继续傻等。这为实现超时获取锁、取消阻塞操作提供了基础。
public void lockInterruptibly() throws InterruptedException {
sync.acquireInterruptibly(1);
}
// AQS.acquireInterruptibly()
public final void acquireInterruptibly(int arg) throws InterruptedException {
if (Thread.interrupted())
throw new InterruptedException();
if (!tryAcquire(arg))
doAcquireInterruptibly(arg); // 可中断的等待逻辑
}
tryLock(long timeout, TimeUnit unit)允许在指定时间内尝试获取锁,超时直接返回false。它的底层基于LockSupport.parkNanos()进行限时等待,同时会检测中断和超时两种信号。
// 超时获取锁使用示例
public boolean tryLockWithTimeout(ReentrantLock lock, long timeoutMs) {
try {
return lock.tryLock(timeoutMs, TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return false;
}
}
这个特性在分布式系统中尤其重要——能有效避免因网络波动导致的线程永久阻塞。
ReentrantLock默认采用非公平锁,但可以通过构造函数参数指定为公平锁:
ReentrantLock fairLock = new ReentrantLock(true); // 公平锁 ReentrantLock unfairLock = new ReentrantLock(false); // 非公平锁(默认)
公平锁严格按照线程请求的顺序来分配锁,新线程来了必须乖乖排到队尾;非公平锁则允许新线程在锁释放的那一刻直接入场“抢跑”,不管队列里还有没有人在等。
公平锁的tryAcquire()会额外检查hasQueuedPredecessors(),确保只有等待队列中没有前驱节点时才尝试获取锁。
从性能角度看,非公平锁通常更高(吞吐量能高出约30%),但代价是有可能导致线程饥饿。公平锁保证了公平,但增加了上下文切换的开销。
/**
* ReentrantLock完整使用示例
*/
public class ReentrantLockDemo {
private final ReentrantLock lock = new ReentrantLock(true);
private final Condition notFull = lock.newCondition();
private final Condition notEmpty = lock.newCondition();
private final Queue queue = new LinkedList<>();
private static final int CAPACITY = 10;
public void produce(int data) throws InterruptedException {
lock.lock();
try {
while (queue.size() == CAPACITY) {
notFull.await(); // 队列满,生产者等待
}
queue.offer(data);
notEmpty.signal(); // 唤醒消费者
} finally {
lock.unlock(); // 必须在finally中释放锁
}
}
public int consume() throws InterruptedException {
lock.lockInterruptibly(); // 可中断方式获取锁
try {
while (queue.isEmpty()) {
notEmpty.await();
}
int data = queue.poll();
notFull.signal();
return data;
} finally {
lock.unlock();
}
}
}
几个关键注意点:
unlock()必须放在finally块里,否则一旦出异常,锁就容易“泄漏”
多个Condition可以精确控制唤醒条件,这是synchronized做不到的
JDK7的ConcurrentHashMap采用\ 分段锁(Segment)\ 设计,整个哈希表被拆成16个Segment数组,每个Segment独立加锁。
核心数据结构长这样:
ConcurrentHashMap
└── Segment[] (默认16个)
└── HashEntry[] (每个Segment独立的哈希表)
└── HashEntry链表
Segment继承了ReentrantLock,每次put操作只锁定对应的Segment,其他Segment的读写不受影响。理论上,最大并发度等于Segment的数量(默认16)。
但分段锁有几个明显的缺陷:
并发度固定,数组扩容后并发能力并不会跟着提升
每个Segment都需要独立的锁和数据结构,浪费内存
跨段操作(比如size())需要锁定所有Segment,性能一下就下来了
JDK8彻底放弃了分段锁,改用数组+链表+红黑树的结构,和HashMap保持了一致。锁的粒度也从Segment级别细化到了每个哈希桶(Node)级别。
核心改进点:
用synchronized替换了ReentrantLock,因为JVM对synchronized的优化更成熟
无锁竞争时,直接用CAS做无锁化更新
只有在哈希桶发生冲突时,才对桶的首节点加锁
链表长度超过8时自动转成红黑树,解决哈希碰撞攻击
// JDK8 ConcurrentHashMap.putVal()核心逻辑
final V putVal(K key, V value, boolean onlyIfAbsent) {
int hash = spread(key.hashCode());
for (Node[] tab = table;;) {
int i = (n - 1) & hash;
Node f = tabAt(tab, i);
if (f == null) {
// 桶为空,CAS直接插入
if (casTabAt(tab, i, null, new Node(hash, key, value)))
break;
} else {
synchronized (f) { // 仅锁定桶的首节点
if (tabAt(tab, i) == f) { // 双重检查
// 链表或红黑树插入逻辑
}
}
}
}
}
JDK8放弃分段锁,背后有四个关键考量:
1. 锁粒度更细
分段锁最小粒度是Segment,JDK8最小粒度是哈希桶
并发度随数组容量动态提升,理论最大并发度等于数组长度
2. synchronized 性能优化
JDK6之后synchronized引入了偏向锁、轻量级锁、自适应自旋
无锁竞争时,偏向锁的性能甚至优于ReentrantLock
JVM可以对synchronized做逃逸分析等深度优化
3. 减少内存开销
消除了Segment对象的内存占用
数据结构与HashMap统一,代码复用性更高
4. 红黑树引入
解决了哈希碰撞导致链表过长的问题
极端情况下,时间复杂度从O(n)降到了O(logn)
ConcurrentHashMap的扩容是并发协作式的,支持多线程一起参与搬砖,这是它的一个亮点。
扩容触发条件:
元素数量达到阈值(容量×加载因子)
单桶链表长度超过8但数组容量小于64
扩容核心流程:
扩容准备:创建nextTable,容量是原数组的2倍
扩容标记:sizeCtl设为负数,表明正在扩容中
任务分配:每个线程负责连续的16个桶的迁移
并发迁移:
处理完的桶会被标记为ForwardingNode
遇到ForwardingNode自动跳过,或者主动协助扩容
扩容完成:table指向nextTable,重置sizeCtl
// 扩容时的ForwardingNode标记 static final class ForwardingNodeextends Node { final Node [] nextTable; ForwardingNode(Node [] tab) { super(MOVED, null, null, null); this.nextTable = tab; } }
并发扩容的巧妙之处:
读操作遇到ForwardingNode会直接转发到新数组,不阻塞
写操作遇到ForwardingNode会主动帮忙扩容,实现了“多线程搭把手”的效果
迁移采用“复制+清除”方式,保证数据一致性
以put操作为例,整个流程是这样的:
计算key的哈希值,定位哈希桶位置
如果数组没初始化,用CAS触发初始化
如果目标桶是空的,CAS直接插入新节点
如果遇到ForwardingNode,帮忙扩容后重试
否则,对桶的首节点加synchronized锁
遍历链表或红黑树,key存在就更新,不存在就追加
链表长度超过8时,触发树化或扩容
释放锁,CAS更新元素计数,检查是否需要扩容
整个过程中,锁的持有时间非常短,绝大多数操作都是无锁的CAS——这正是ConcurrentHashMap高性能的根本原因。
BlockingQueue是Ja va并发包里最重要的数据结构之一,专门用来解决生产者-消费者模式下的线程协作问题。
核心特性:
队列满时,生产者线程自动阻塞,直到有消费者来消费
队列空时,消费者线程自动阻塞,直到有生产者生产
所有操作都是线程安全的,内部通过锁和条件变量实现
四种核心操作模式:
操作方式 | 抛出异常 | 返回特殊值 | 阻塞等待 | 超时等待 |
插入 | add(e) | offer(e) | put(e) | offer(e, time, unit) |
移除 | remove() | poll() | take() | poll(time, unit) |
检查 | element() | peek() | - | - |
可以说,阻塞队列不仅是线程池的核心组件,也是解耦生产者和消费者速率不匹配的标准方案。
ArrayBlockingQueue是基于数组实现的有界阻塞队列,创建时必须指定容量。
核心特点:
有界队列,容量固定,不能扩容
单锁双Condition机制(notEmpty + notFull)
支持公平/非公平模式
读写共用同一把锁,并发度相对较低
// ArrayBlockingQueue核心结构 public class ArrayBlockingQueue{ final Object[] items; // 存储数组 int takeIndex; // 取元素位置 int putIndex; // 放元素位置 int count; // 元素数量 final ReentrantLock lock; // 单锁 private final Condition notEmpty; private final Condition notFull; }
适用场景:队列大小可以预估、对内存占用比较敏感的场景。
LinkedBlockingQueue是基于链表实现的阻塞队列,默认容量为Integer.MAX_VALUE。
核心特点:
双锁分离设计(takeLock + putLock),读写互不阻塞
默认无界(实际最大2^31-1),也可以指定容量
吞吐量高于ArrayBlockingQueue
内存占用相对较高
// LinkedBlockingQueue双锁设计 private final ReentrantLock takeLock = new ReentrantLock(); private final Condition notEmpty = takeLock.newCondition(); private final ReentrantLock putLock = new ReentrantLock(); private final Condition notFull = putLock.newCondition();
需要留意:在无界模式下,如果生产速度远大于消费速度,OOM可不是开玩笑的。
适用场景:生产消费速率差异较大、追求高吞吐量的场景。
SynchronousQueue是一个不存储元素的阻塞队列——每个插入操作必须等待另一个线程的移除操作。
核心特点:
内部没有缓冲区,队列容量始终为0
支持公平/非公平模式
直接传递,不存储任何元素
是Executors.newCachedThreadPool()的默认队列
// SynchronousQueue典型用法 SynchronousQueuequeue = new SynchronousQueue<>(); // 生产者线程 new Thread(() -> { queue.put(1); // 会阻塞直到有消费者take }).start(); // 消费者线程 new Thread(() -> { queue.take(); // 会阻塞直到有生产者put }).start();
适用场景:任务必须立即处理、不允许排队的场景,比如CachedThreadPool。
DelayQueue是支持延时获取元素的无界阻塞队列,元素必须实现Delayed接口。
核心原理:
内部基于PriorityQueue实现,按过期时间排序
只有元素延迟时间到期后才能被取出
队首永远是最早过期的元素
// 延时任务元素定义
public class DelayedTask implements Delayed {
private final long executeTime;
private final Runnable task;
public DelayedTask(long delayMs, Runnable task) {
this.executeTime = System.currentTimeMillis() + delayMs;
this.task = task;
}
@Override
public long getDelay(TimeUnit unit) {
return unit.convert(executeTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS);
}
@Override
public int compareTo(Delayed other) {
return Long.compare(this.executeTime, ((DelayedTask)other).executeTime);
}
}
典型应用场景:
订单超时自动取消
会话超时清理
定时任务调度
缓存过期失效
// 延时队列使用示例 DelayQueuedelayQueue = new DelayQueue<>(); delayQueue.put(new DelayedTask(30000, () -> System.out.println("30秒后执行"))); delayQueue.put(new DelayedTask(60000, () -> System.out.println("60秒后执行"))); // 消费者线程 while (true) { DelayedTask task = delayQueue.take(); // 阻塞直到任务到期 task.run(); }
特性 | ArrayBlockingQueue | LinkedBlockingQueue | SynchronousQueue | DelayQueue |
容量 | 有界 | 可选有界 / 无界 | 0 | 无界 |
数据结构 | 数组 | 链表 | 直接传递 | 优先级队列 |
锁机制 | 单锁 | 双锁分离 | CAS | 单锁 |
公平性 | 支持 | 不支持 | 支持 | 不支持 |
吞吐量 | 中等 | 高 | 极高 | 低 |
典型应用 | 固定大小线程池 | 固定大小线程池 | 缓存线程池 | 延时任务 |
选型建议:
需要控制队列大小 → ArrayBlockingQueue
追求高吞吐量 → LinkedBlockingQueue
任务必须立即执行 → SynchronousQueue
需要延时执行 → DelayQueue
传统Future接口的局限性其实很明显:
无法手动完成
不支持链式调用
不支持异常处理
无法组合多个Future
CompletableFuture在Ja va8中被引入,实现了CompletionStage接口,提供了丰富的异步编程能力,而且支持函数式编程风格,用起来非常灵活。
1. 创建类 API
// 使用默认线程池 CompletableFuturefuture1 = CompletableFuture.runAsync(() -> System.out.println("异步任务")); CompletableFuture future2 = CompletableFuture.supplyAsync(() -> "返回结果"); // 使用自定义线程池(推荐) ExecutorService executor = Executors.newFixedThreadPool(10); CompletableFuture future3 = CompletableFuture.supplyAsync(() -> "自定义线程池", executor); // 手动完成 CompletableFuture future4 = new CompletableFuture<>(); future4.complete("手动设置结果"); future4.completeExceptionally(new RuntimeException("手动异常"));
需要提醒的是:默认使用ForkJoinPool.commonPool(),所有CompletableFuture共享这个线程池。如果是CPU密集型任务,强烈建议用自定义线程池。
2. 链式转换类 API
CompletableFuturefuture = CompletableFuture.supplyAsync(() -> "Hello") .thenApply(s -> s + " World") // 同步转换 .thenApplyAsync(s -> s.toUpperCase()) // 异步转换 .thenAccept(s -> System.out.println(s)) // 消费结果 .thenRun(() -> System.out.println("执行完成")); // 仅执行,不消费结果
thenApply:输入T,输出钱,类似map操作
thenAccept:输入T,无输出,消费型
thenRun:无输入无输出,只执行动作
3. 组合类 API
// AND组合:两个都完成才执行 CompletableFuturef1 = CompletableFuture.supplyAsync(() -> "Hello"); CompletableFuture f2 = CompletableFuture.supplyAsync(() -> "World"); f1.thenCombine(f2, (s1, s2) -> s1 + " " + s2) .thenAccept(System.out::println); // 输出 Hello World // OR组合:任意一个完成就执行 CompletableFuture fast = f1.applyToEither(f2, s -> s + " faster");
4. 异常处理 API
CompletableFuture.supplyAsync(() -> {
if (true) throw new RuntimeException("出错了");
return "正常";
})
.exceptionally(ex -> {
System.out.println("捕获异常: " + ex.getMessage());
return "默认值"; // 异常时返回默认值
})
.handle((result, ex) -> {
if (ex != null) {
return "处理异常";
}
return result;
})
.whenComplete((result, ex) -> {
// 无论成功失败都会执行,不改变结果
System.out.println("执行完成");
});
5. 多任务组合
// 所有任务都完成 CompletableFutureall = CompletableFuture.allOf(f1, f2, f3); // 任意一个任务完成 CompletableFuture
/**
* 并行调用多个服务,聚合结果
*/
public class CompletableFutureDemo {
public UserInfo getUserInfo(Long userId) {
// 并行调用三个接口
CompletableFuture userFuture = CompletableFuture.supplyAsync(() -> userService.getUser(userId));
CompletableFuture> orderFuture = CompletableFuture.supplyAsync(() -> orderService.getOrders(userId));
CompletableFuture> addrFuture = CompletableFuture.supplyAsync(() -> addressService.getAddresses(userId));
// 等待所有完成,聚合结果
return CompletableFuture.allOf(userFuture, orderFuture, addrFuture)
.thenApply(v -> {
UserInfo info = new UserInfo();
info.setUser(userFuture.join());
info.setOrders(orderFuture.join());
info.setAddresses(addrFuture.join());
return info;
})
.exceptionally(ex -> {
log.error("获取用户信息失败", ex);
return null;
})
.join();
}
}
通过CompletableFuture,原本串行需要3秒的调用,可以优化到1秒内完成——这在微服务架构下是非常常用的优化手段。
这篇文章系统梳理了Ja va并发编程的四大核心模块:ReentrantLock的灵活锁机制、ConcurrentHashMap的高性能并发设计、BlockingQueue的生产者-消费者模式,以及CompletableFuture的异步编程能力。
这些技术是构建高并发系统的基石。建议结合具体的项目去深入实践,重点关注各组件的适用场景和性能特性,避免在生产环境中间出现并发安全问题。如果想继续深入,可以去研究AQS框架、JMM内存模型以及无锁算法,这些才是真正让你对并发编程“豁然开朗”的东西。
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
正版软件
正版软件
正版软件
正版软件
正版软件
1
2
3
7
8