JUC 并发编程
概述
JUC 是 java.util.concurrent 包的缩写,Java 5 引入,提供了丰富的并发编程工具。
| 组件 | 说明 |
|---|---|
| Lock 框架 | ReentrantLock、ReadWriteLock、StampedLock |
| AQS | AbstractQueuedSynchronizer,大部分同步器的基石 |
| 并发集合 | ConcurrentHashMap、CopyOnWriteArrayList、BlockingQueue 等 |
| 线程池 | ThreadPoolExecutor、ScheduledThreadPoolExecutor |
| 并发工具类 | CountDownLatch、CyclicBarrier、Semaphore、Exchanger |
| Fork/Join | 分治任务的并行执行框架 |
| CompletableFuture | 异步编程的利器 |
Lock 接口与 ReentrantLock
Lock vs synchronized
| 对比维度 | synchronized | Lock |
|---|---|---|
| 用法 | JVM 关键字,自动加解锁 | API 接口,需手动 lock/unlock |
| 锁释放 | 退出同步块自动释放 | 必须在 finally 中 unlock |
| 中断响应 | 不支持 | 支持 lockInterruptibly() |
| 超时 | 不支持 | 支持 tryLock(timeout) |
| 公平性 | 非公平 | 可设置公平/非公平 |
| 条件变量 | wait/notify | Condition.await/signal |
ReentrantLock 示例
java
public class Counter {
private final Lock lock = new ReentrantLock();
private int count;
public void increment() {
lock.lock();
try {
count++;
} finally {
lock.unlock(); // 必须释放
}
}
public boolean tryIncrement(long timeout, TimeUnit unit) {
try {
if (lock.tryLock(timeout, unit)) {
try {
count++;
return true;
} finally {
lock.unlock();
}
}
return false;
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return false;
}
}
}公平锁 vs 非公平锁
java
ReentrantLock fairLock = new ReentrantLock(true); // 公平锁
ReentrantLock unfairLock = new ReentrantLock(false); // 非公平锁(默认)- 非公平锁:允许插队,吞吐量更高
- 公平锁:先到先得,避免线程饥饿
ReadWriteLock
读写分离,读读不互斥,读写/写写互斥:
java
public class Cache {
private final Map<String, Object> map = new HashMap<>();
private final ReadWriteLock rw = new ReentrantReadWriteLock();
public Object get(String key) {
rw.readLock().lock();
try { return map.get(key); }
finally { rw.readLock().unlock(); }
}
public void put(String key, Object value) {
rw.writeLock().lock();
try { map.put(key, value); }
finally { rw.writeLock().unlock(); }
}
}StampedLock
Java 8 引入,支持乐观读,在读多写少场景下性能更高:
java
public class Point {
private double x, y;
private final StampedLock sl = new StampedLock();
public double distanceFromOrigin() {
long stamp = sl.tryOptimisticRead(); // 乐观读
double curX = x, curY = y;
if (!sl.validate(stamp)) { // 被写过?
stamp = sl.readLock(); // 升级为悲观读
try {
curX = x;
curY = y;
} finally {
sl.unlockRead(stamp);
}
}
return Math.sqrt(curX * curX + curY * curY);
}
}AQS 原理
AbstractQueuedSynchronizer 是 JUC 的基石,ReentrantLock、CountDownLatch、Semaphore 等都基于它实现。
核心思想
volatile int state ← 同步状态(0=未锁,1=锁定,N=重入)
│
┌──────────┴──────────┐
│ │
独占模式 共享模式
(ReentrantLock) (Semaphore/CountDownLatch)
│ │
└──────────┬──────────┘
│
CLH 队列(双向链表)
┌─► Node ◄─► Node ◄─► Node ─┐
│ 等待线程 等待线程 等待线程 │
└───────────────────────────────┘核心方法
| 方法 | 说明 |
|---|---|
tryAcquire(int) | 独占式获取状态 |
tryRelease(int) | 独占式释放状态 |
tryAcquireShared(int) | 共享式获取状态 |
tryReleaseShared(int) | 共享式释放状态 |
isHeldExclusively() | 是否独占模式 |
自定义同步器示例
java
public class SharedLock {
private static class Sync extends AbstractQueuedSynchronizer {
Sync(int count) {
setState(count);
}
@Override
protected int tryAcquireShared(int acquires) {
while (true) {
int available = getState();
int remaining = available - acquires;
if (remaining < 0 || compareAndSetState(available, remaining))
return remaining;
}
}
@Override
protected boolean tryReleaseShared(int releases) {
while (true) {
int current = getState();
int next = current + releases;
if (compareAndSetState(current, next))
return true;
}
}
}
private final Sync sync = new Sync(3); // 最多 3 个并发
public void acquire() { sync.acquireShared(1); }
public void release() { sync.releaseShared(1); }
}并发集合
ConcurrentHashMap
| 版本 | 实现方式 |
|---|---|
| Java 7 | Segment 分段锁 |
| Java 8+ | Node + CAS + synchronized(桶内链表/红黑树) |
常用方法:
java
ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<>();
map.putIfAbsent("key", 1); // 不存在才放入
map.computeIfAbsent("key", k -> expensiveLoad(k)); // 原子计算
map.merge("key", 1, Integer::sum); // 原子合并
map.forEach(1, (k, v) -> System.out.println(k + "=" + v)); // 并行遍历CopyOnWriteArrayList
读操作不加锁,写操作复制新数组,适合读多写极少场景:
java
CopyOnWriteArrayList<String> list = new CopyOnWriteArrayList<>();
list.add("A"); // 内部复制数组,加锁
list.get(0); // 无锁,直接读取缺点:写操作开销大,数据最终一致性。
BlockingQueue
| 实现类 | 特点 |
|---|---|
ArrayBlockingQueue | 有界数组,公平模式可选 |
LinkedBlockingQueue | 可选有界,链表结构 |
PriorityBlockingQueue | 无界,支持优先级 |
DelayQueue | 延迟出队 |
SynchronousQueue | 容量为 0,直接传递 |
LinkedTransferQueue | 支持 transfer 方法 |
方法对比(4 种处理方式):
| 抛异常 | 返回特殊值 | 阻塞 | 超时 | |
|---|---|---|---|---|
| 入队 | add(e) | offer(e) | put(e) | offer(e, t, u) |
| 出队 | remove() | poll() | take() | poll(t, u) |
线程池
ThreadPoolExecutor 构造
java
public ThreadPoolExecutor(
int corePoolSize, // 核心线程数
int maximumPoolSize, // 最大线程数
long keepAliveTime, // 空闲线程存活时间
TimeUnit unit, // 时间单位
BlockingQueue<Runnable> workQueue, // 任务队列
ThreadFactory threadFactory, // 线程工厂
RejectedExecutionHandler handler // 拒绝策略
);运行流程
提交任务
│
▼
当前线程 < corePoolSize? ──是──► 新建核心线程执行
│ 否
▼
任务入队成功? ───是──► 等待核心线程处理
│ 否
▼
当前线程 < maximumPoolSize? ──是──► 新建非核心线程执行
│ 否
▼
执行拒绝策略拒绝策略
| 策略 | 行为 |
|---|---|
AbortPolicy(默认) | 抛 RejectedExecutionException |
CallerRunsPolicy | 调用者线程执行 |
DiscardPolicy | 直接丢弃 |
DiscardOldestPolicy | 丢弃队首,重新提交 |
Executors 工厂
java
// 固定大小线程池(LinkedBlockingQueue,可能 OOM)
ExecutorService fixed = Executors.newFixedThreadPool(4);
// 缓存线程池(SynchronousQueue,最大为 Integer.MAX_VALUE)
ExecutorService cached = Executors.newCachedThreadPool();
// 单线程线程池
ExecutorService single = Executors.newSingleThreadExecutor();
// 定时任务
ScheduledExecutorService scheduled = Executors.newScheduledThreadPool(2);建议:手动创建 ThreadPoolExecutor,明确参数。
正确关闭线程池
java
executor.shutdown(); // 不再接受新任务
try {
if (!executor.awaitTermination(5, TimeUnit.SECONDS)) {
executor.shutdownNow(); // 强制停止
if (!executor.awaitTermination(5, TimeUnit.SECONDS)) {
System.err.println("线程池未终止");
}
}
} catch (InterruptedException e) {
executor.shutdownNow();
Thread.currentThread().interrupt();
}CompletableFuture
Java 8 引入的异步编程工具,支持函数式组合。
基础用法
java
// 异步执行
CompletableFuture.supplyAsync(() -> fetchData())
.thenApply(data -> process(data))
.thenAccept(result -> save(result))
.exceptionally(e -> {
log.error("处理失败", e);
return null;
});组合方式
java
// 两个独立任务并行
CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> task1());
CompletableFuture<String> f2 = CompletableFuture.supplyAsync(() -> task2());
f1.thenCombine(f2, (r1, r2) -> r1 + r2); // 两个都完成
f1.acceptEither(f2, r -> log.info(r)); // 任一完成
CompletableFuture.anyOf(f1, f2); // 任一完成
CompletableFuture.allOf(f1, f2); // 全部完成
// 嵌套扁平化
CompletableFuture.supplyAsync(() -> task())
.thenCompose(result -> CompletableFuture.supplyAsync(() -> next(result)));自定义线程池
java
ExecutorService executor = Executors.newFixedThreadPool(4);
CompletableFuture.supplyAsync(() -> fetchData(), executor);超时控制(Java 9+)
java
CompletableFuture.supplyAsync(() -> fetchData())
.completeOnTimeout(defaultValue, 3, TimeUnit.SECONDS) // 超时返回默认值
.orTimeout(5, TimeUnit.SECONDS); // 超时抛 TimeoutException并发工具类
CountDownLatch
允许一个或多个线程等待其他线程完成操作:
java
CountDownLatch latch = new CountDownLatch(3);
// 工作线程
void worker() {
doWork();
latch.countDown(); // 计数器 -1
}
// 等待线程
latch.await(); // 阻塞直到计数器归零特点:计数器不可重置。
CyclicBarrier
让一组线程互相等待到达某个屏障,然后继续:
java
CyclicBarrier barrier = new CyclicBarrier(3, () -> {
System.out.println("所有线程到达屏障,开始下一阶段");
});
void worker() {
prepare();
barrier.await(); // 等待其他线程
execute();
}与 CountDownLatch 对比:
| 特性 | CountDownLatch | CyclicBarrier |
|---|---|---|
| 重置 | 不可重置 | 可 reset 重用 |
| 角色 | 一个等 N 个 | N 个互相等 |
| 屏障动作 | 无 | 支持回调 |
Semaphore
控制同时访问的线程数:
java
Semaphore sem = new Semaphore(3); // 最多 3 个
void access() {
try {
sem.acquire(); // 获取许可
doAccess();
} finally {
sem.release(); // 释放许可
}
}Exchanger
两个线程交换数据:
java
Exchanger<String> exchanger = new Exchanger<>();
// 线程 A
String dataFromB = exchanger.exchange("来自A的数据");
// 线程 B
String dataFromA = exchanger.exchange("来自B的数据");线程安全最佳实践
- 优先使用并发容器,而非手动加锁
- 线程池统一管理,避免
new Thread() - 锁的粒度要小,避免在锁内做 IO 操作
- 注意可见性:final、volatile、锁
- 避免死锁:按固定顺序获取锁,使用
tryLock - 线程命名:设置有意义的线程名方便排错
- 慎用全局锁,优先局部锁或 CAS
- 异步异常处理:CompletableFuture 务必接
exceptionally