Java 并发入门:线程、锁与线程池全景
体系化梳理 Java 并发编程的核心工具与适用场景:从线程基础、三大并发问题,到 synchronized / Lock / 原子类 / 线程池 / 并发工具类,每个工具都配有典型场景与 JDK 8 兼容代码示例,帮助快速建立并发编程全局视图。
目录
| 章节 | 说明 |
|---|---|
| 为什么需要并发 | 并发的价值与代价 |
| 线程基础 | 创建线程、状态机、基本操作 |
| 三大并发问题 | 可见性、原子性、有序性 |
| 同步工具 | synchronized、volatile、Lock、ReadWriteLock |
| 并发工具类 | CountDownLatch、CyclicBarrier、Semaphore |
| 线程安全集合 | ConcurrentHashMap、阻塞队列、CopyOnWriteArrayList |
| 原子类 | AtomicXxx、LongAdder、ABA 问题 |
| 线程池 | ThreadPoolExecutor 核心参数与实战 |
| 快速选型指南 | 场景 → 工具的决策树 |
为什么需要并发
并发的价值:
- 提升吞吐量:IO 密集型任务(读数据库、调接口)等待期间,让 CPU 处理其他请求
- 充分利用多核:CPU 密集型任务可并行分治,用多核加速
- 提升响应速度:后台异步任务不阻塞主流程,用户感知更快
并发的代价:引入了可见性、原子性、有序性三类线程安全问题(详见下文)。
选择并发的判断标准:任务之间有没有依赖关系?相互独立 → 可并发;有依赖 → 需要同步协调。
线程基础
创建线程的三种方式
import java.util.concurrent.*;
// 方式一:继承 Thread(不推荐,占用 Java 的单继承槽位)
class PrintThread extends Thread {
@Override
public void run() {
System.out.println("Thread: " + Thread.currentThread().getName());
}
}
new PrintThread().start();
// 方式二:实现 Runnable(推荐,解耦任务逻辑与线程管理)
// 匿名内部类写法(JDK 7-)
Thread t1 = new Thread(new Runnable() {
@Override
public void run() {
System.out.println("Runnable task");
}
}, "my-thread");
// Lambda 写法(JDK 8+)
Thread t2 = new Thread(() -> System.out.println("Lambda task"), "lambda-thread");
t2.start();
// 方式三:Callable + Future(有返回值,异常可向调用方传递)
ExecutorService executor = Executors.newFixedThreadPool(4);
// Lambda 写法(编译器推断为 Callable<Integer>,允许 checked exception)
Future<Integer> future = executor.submit(() -> {
Thread.sleep(100); // 模拟耗时计算
return 42;
});
try {
System.out.println(future.get()); // 阻塞等待结果
} catch (ExecutionException e) {
System.err.println("任务失败:" + e.getCause().getMessage());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
executor.shutdown();
线程状态机
flowchart LR
NEW -->|"start()"| RUNNABLE
RUNNABLE -->|"竞争 synchronized 锁失败"| BLOCKED
BLOCKED -->|"获得锁"| RUNNABLE
RUNNABLE -->|"Object.wait()"| WAITING
WAITING -->|"notify() / notifyAll()"| RUNNABLE
RUNNABLE -->|"sleep(n) / join(n)"| TIMED_WAITING
TIMED_WAITING -->|"超时 / 唤醒"| RUNNABLE
RUNNABLE -->|"run() 结束"| TERMINATED
style NEW fill:#e8f4f8,stroke:#4a90d9
style TERMINATED fill:#f8e8e8,stroke:#c00
style RUNNABLE fill:#e8f8e8,stroke:#060
| 状态 | 含义 | 进入原因 |
|---|---|---|
NEW | 已创建,未启动 | new Thread() 之后,start() 之前 |
RUNNABLE | 运行中或就绪(等 CPU 调度) | start() 之后 |
BLOCKED | 等待 synchronized 锁 | 竞争 synchronized 代码块 |
WAITING | 无限期等待 | Object.wait()、Thread.join() |
TIMED_WAITING | 有限期等待 | Thread.sleep(n)、Object.wait(n) |
TERMINATED | 已结束 | run() 正常返回或抛出异常 |
常用线程操作
Thread t = new Thread(() -> { /* 任务逻辑 */ });
t.start(); // 启动线程,JVM 调度后执行 run()
t.join(); // 当前线程阻塞,等待 t 执行完毕
t.join(1000); // 最多等 1 秒,不保证 t 已结束
t.interrupt(); // 发送中断信号(只设标志位,不强制停止线程)
t.isInterrupted(); // 检查中断标志(不清除)
Thread.interrupted(); // 检查并清除当前线程的中断标志
Thread.sleep(100); // 当前线程休眠 100ms,释放 CPU 但不释放锁
Thread.currentThread(); // 获取当前线程对象
⚠️ 不要使用
Thread.stop()(已废弃):强制停止可能导致对象处于不一致状态。正确做法是用interrupt()+ 主动检查标志位来协作退出。
三大并发问题
可见性:CPU 缓存的代价
多核 CPU 每个核心有独立的 L1/L2 缓存,线程 A 修改变量后写入缓存,线程 B 从另一个核心的缓存读取时可能看到旧值。
// ❌ 可能死循环:reader 线程可能永远看不到 running = false
class StopThreadBad {
private static boolean running = true;
public static void main(String[] args) throws Exception {
new Thread(new Runnable() {
@Override
public void run() {
while (running) { /* 持续循环 */ }
System.out.println("线程退出"); // 可能永远不会执行
}
}).start();
Thread.sleep(1000);
running = false; // 对 reader 线程不可见!
}
}
// ✅ 修复:加 volatile,强制读写直接访问主内存
class StopThreadOk {
private static volatile boolean running = true;
// ... 其他代码不变
}
原子性:线程切换的代价
count++ 在 JVM 层面是三条指令:读内存 → 加 1 → 写内存。OS 可以在任意两条指令之间切换线程,导致两个线程基于同一旧值计算,相互覆盖写入结果。
// ❌ 两个线程各执行 10000 次,最终结果不是 20000
class UnsafeCounter {
long count = 0;
void add() { count++; } // 读-改-写三步,非原子
}
// ✅ 方案一:synchronized 加锁(保证原子性 + 可见性)
class LockedCounter {
long count = 0;
synchronized void add() { count++; }
synchronized long get() { return count; }
}
// ✅ 方案二:AtomicLong(无锁,性能更好)
class AtomicCounter {
AtomicLong count = new AtomicLong(0);
void add() { count.incrementAndGet(); }
long get() { return count.get(); }
}
有序性:编译优化的代价
编译器和 CPU 为提升性能会对指令重排序。最经典的问题是双重检查单例(DCL):
// ❌ 问题:instance 未加 volatile,new Singleton() 的步骤可能被重排序
// 正常步骤:① 分配内存 → ② 初始化对象 → ③ 赋值给 instance
// 重排后可能:① 分配内存 → ② 赋值给 instance(非空!)→ ③ 初始化对象
// 线程 B 看到 instance 非空就直接返回,但对象还没初始化完!
class SingletonBad {
static SingletonBad instance;
static SingletonBad getInstance() {
if (instance == null) {
synchronized (SingletonBad.class) {
if (instance == null)
instance = new SingletonBad(); // 危险!可能重排序
}
}
return instance;
}
}
// ✅ 修复:volatile 禁止 new 操作的指令重排序
class SingletonOk {
static volatile SingletonOk instance;
static SingletonOk getInstance() {
if (instance == null) {
synchronized (SingletonOk.class) {
if (instance == null)
instance = new SingletonOk(); // volatile 保证赋值在初始化之后
}
}
return instance;
}
}
同步工具
synchronized — 最简单的互斥锁
适用场景:保护临界区(需要原子执行的代码块),竞争不激烈时首选,代码简洁,自动释放锁。
class BankAccount {
private long balance;
// 场景 1:方法级锁,锁对象 = this(实例方法)
public synchronized void deposit(long amount) {
balance += amount;
}
// 场景 2:代码块锁,锁粒度更小,减少持锁时间
public void withdraw(long amount) {
if (amount <= 0) throw new IllegalArgumentException("金额非法");
synchronized (this) { // 只锁核心操作
if (balance < amount) throw new IllegalStateException("余额不足");
balance -= amount;
}
}
// 场景 3:静态方法锁,锁对象 = BankAccount.class
public static synchronized BankAccount createDefault() {
return new BankAccount();
}
public synchronized long getBalance() { return balance; }
}
⚠️ 常见错误:
// ❌ 每次 new 一个新锁对象,完全无保护效果
synchronized (new Object()) { count++; }
// ❌ add() 和 sub() 用了不同的锁,两个方法可以同时执行
class Wrong {
private final Object lock = new Object();
long count = 0;
synchronized void add() { count++; } // 锁 this
void sub() { synchronized (lock) { count--; } } // 锁 lock(不同对象!)
}
volatile — 轻量级可见性与有序性保证
适用场景:状态标志位、单次写-多次读的场景。不保证原子性,不能替代 synchronized。
// 场景 1:停止后台线程
class BackgroundWorker {
private volatile boolean shutdown = false;
public void run() {
while (!shutdown) {
doWork();
}
System.out.println("Worker 已停止");
}
public void stop() {
shutdown = true; // 立即对其他线程可见
}
}
// 场景 2:DCL 单例(见三大并发问题示例)
// ❌ volatile 不保证原子性!
class NotSafe {
private volatile long count = 0;
void add() { count++; } // count++ 仍是三步操作,线程不安全
}
Lock — 灵活的显式锁
适用场景:需要以下特性时替代 synchronized:
- 超时获锁(
tryLock(timeout)):避免死等,可降级处理 - 可中断获锁(
lockInterruptibly()):支持任务取消 - 多个条件变量(
Condition):精确唤醒某类等待线程
import java.util.concurrent.locks.*;
import java.util.LinkedList;
import java.util.Queue;
// 场景:有界阻塞缓冲区,生产者/消费者精确唤醒
class BoundedBuffer<T> {
private final Lock lock = new ReentrantLock();
private final Condition notFull = lock.newCondition(); // "未满"条件
private final Condition notEmpty = lock.newCondition(); // "非空"条件
private final Queue<T> queue = new LinkedList<T>();
private final int capacity;
public BoundedBuffer(int capacity) { this.capacity = capacity; }
// 生产者:队列满时阻塞,等待"未满"信号
public void put(T item) throws InterruptedException {
lock.lock();
try {
while (queue.size() == capacity) {
notFull.await(); // 释放锁并等待
}
queue.add(item);
notEmpty.signal(); // 通知一个消费者
} finally {
lock.unlock(); // 必须在 finally 中释放,防止异常时死锁
}
}
// 消费者:队列空时阻塞,等待"非空"信号
public T take() throws InterruptedException {
lock.lock();
try {
while (queue.isEmpty()) {
notEmpty.await();
}
T item = queue.poll();
notFull.signal();
return item;
} finally {
lock.unlock();
}
}
}
// 场景:超时获锁,避免死等
class TimeoutExample {
private final ReentrantLock lock = new ReentrantLock();
void tryUpdate() throws InterruptedException {
if (lock.tryLock(500, java.util.concurrent.TimeUnit.MILLISECONDS)) {
try {
// 成功获锁,执行操作
doUpdate();
} finally {
lock.unlock();
}
} else {
// 超时,执行降级逻辑
System.out.println("获锁超时,跳过本次更新");
}
}
private void doUpdate() { /* ... */ }
}
ReadWriteLock — 读多写少场景
读操作并发(多个读者同时持锁),写操作独占(写者持锁时读者阻塞)。
import java.util.concurrent.locks.*;
import java.util.HashMap;
import java.util.Map;
// 场景:本地缓存(读远多于写)
class LocalCache<K, V> {
private final ReadWriteLock rwl = new ReentrantReadWriteLock();
private final Map<K, V> cache = new HashMap<K, V>();
// 多个读线程可并发执行
public V get(K key) {
rwl.readLock().lock();
try {
return cache.get(key);
} finally {
rwl.readLock().unlock();
}
}
// 写操作独占,读写互斥
public void put(K key, V value) {
rwl.writeLock().lock();
try {
cache.put(key, value);
} finally {
rwl.writeLock().unlock();
}
}
}
三种工具对比:
| 工具 | 原子性 | 可见性 | 有序性 | 核心特点 |
|---|---|---|---|---|
synchronized | ✅ | ✅ | ✅ | 简单,自动释放锁 |
volatile | ❌ | ✅ | ✅(禁重排) | 轻量,只保证可见性 |
ReentrantLock | ✅ | ✅ | ✅ | 灵活,需手动释放 |
ReentrantReadWriteLock | ✅ | ✅ | ✅ | 读并发,写独占 |
并发工具类
CountDownLatch — 等待多个任务完成
场景:主线程等待 N 个子任务全部完成后再继续(一次性,不可重置)。
import java.util.concurrent.*;
// 场景:并行初始化多个服务,全部就绪后才开放流量
final int serviceCount = 3;
final CountDownLatch latch = new CountDownLatch(serviceCount);
final String[] services = {"DB", "Cache", "MQ"};
for (final String name : services) {
new Thread(new Runnable() {
@Override
public void run() {
try {
initService(name); // 模拟各服务初始化(耗时不同)
System.out.println(name + " 初始化完成");
} finally {
latch.countDown(); // 无论成功失败,计数减 1
}
}
}).start();
}
latch.await(); // 主线程阻塞,直到计数归零
System.out.println("所有服务就绪,开始接收请求");
// 也可以指定最大等待时间
boolean done = latch.await(30, TimeUnit.SECONDS);
if (!done) System.err.println("服务启动超时!");
CyclicBarrier — 多线程同步到达集合点
场景:N 个线程都到达某个同步点后,再一起继续(可重复使用,适合多阶段并行任务)。
import java.util.concurrent.*;
// 场景:分阶段并行计算,每阶段所有线程完成后汇总再进入下一阶段
final int workerCount = 3;
final CyclicBarrier barrier = new CyclicBarrier(workerCount, new Runnable() {
@Override
public void run() {
// 当所有线程都到达 barrier 时,由最后到达的线程执行此回调
System.out.println("=== 阶段完成,开始下一阶段 ===");
}
});
for (int i = 0; i < workerCount; i++) {
final int workerId = i;
new Thread(new Runnable() {
@Override
public void run() {
try {
computePhase1(workerId); // 并行执行第一阶段
barrier.await(); // 等待所有线程完成第一阶段
computePhase2(workerId); // 并行执行第二阶段
barrier.await(); // 等待所有线程完成第二阶段(CyclicBarrier 可重用)
} catch (InterruptedException | BrokenBarrierException e) {
Thread.currentThread().interrupt();
}
}
}).start();
}
Semaphore — 限制并发数量
场景:流量控制(数据库连接池、对外接口限流)。
import java.util.concurrent.*;
// 场景:最多允许 5 个线程同时查询数据库
final Semaphore semaphore = new Semaphore(5);
void queryDatabase(String sql) {
try {
semaphore.acquire(); // 获取 1 个许可(无许可时阻塞)
try {
executeQuery(sql);
} finally {
semaphore.release(); // 必须在 finally 中释放!
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// 也可以非阻塞尝试
if (semaphore.tryAcquire()) {
try {
executeQuery(sql);
} finally {
semaphore.release();
}
} else {
// 无许可可用,快速失败
throw new RuntimeException("服务繁忙,请稍后重试");
}
| 工具 | 核心用途 | 是否可重置 |
|---|---|---|
CountDownLatch | 等待 N 个事件发生 | 否(一次性) |
CyclicBarrier | N 个线程同步到集合点 | 是 |
Semaphore | 控制最大并发数量 | — |
线程安全集合
ConcurrentHashMap
场景:多线程共享 Map,替代线程不安全的 HashMap(易数据丢失/死循环)和性能差的 Hashtable(全表锁)。
import java.util.concurrent.*;
import java.util.function.*;
ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<String, Integer>();
// 普通读写
map.put("key", 1);
Integer val = map.get("key");
// 原子性复合操作(不能用 get + put 代替,那不是原子的!)
map.putIfAbsent("key", 1); // 不存在才插入,返回原来的值(或 null)
// Java 8:compute 系列(lambda 在持锁期间执行,保证原子性)
map.computeIfAbsent("newKey", new Function<String, Integer>() {
@Override
public Integer apply(String k) {
return k.length(); // 不存在时计算并插入
}
});
// Lambda 简写
map.computeIfAbsent("anotherKey", k -> k.length());
// Java 8:merge(存在则合并,适合计数器场景)
map.merge("counter", 1, new BiFunction<Integer, Integer, Integer>() {
@Override
public Integer apply(Integer oldVal, Integer newVal) {
return oldVal + newVal;
}
});
// Lambda 简写
map.merge("counter", 1, Integer::sum); // 不存在则插入 1,存在则累加
⚠️ key 和 value 都不能为 null;并发修改时
size()是近似值。
阻塞队列(生产者-消费者)
场景:解耦生产速度与消费速度,队列满/空时自动阻塞,无需手写 wait/notify。
import java.util.concurrent.*;
// 推荐使用有界队列,防止生产过快导致 OOM
final BlockingQueue<String> queue = new ArrayBlockingQueue<String>(100);
// 生产者线程
Thread producer = new Thread(new Runnable() {
@Override
public void run() {
try {
while (true) {
String task = buildTask();
queue.put(task); // 队列满时阻塞,直到有空间
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
// 消费者线程
Thread consumer = new Thread(new Runnable() {
@Override
public void run() {
try {
while (true) {
String task = queue.take(); // 队列空时阻塞,直到有任务
process(task);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
});
| 实现类 | 是否有界 | 推荐度 | 特点 |
|---|---|---|---|
ArrayBlockingQueue | 是 | ⭐⭐⭐ | 推荐,防 OOM,数组实现 |
LinkedBlockingQueue | 可选(默认无界) | ⭐⭐ | 默认无界时注意 OOM |
SynchronousQueue | 无缓冲 | ⭐⭐ | 生产者必须等消费者接手才返回 |
DelayQueue | 无界 | ⭐⭐ | 按延时时间出队(定时任务) |
PriorityBlockingQueue | 无界 | ⭐⭐ | 按优先级出队(元素需实现 Comparable) |
CopyOnWriteArrayList
场景:读多写极少的 List(如监听器列表、白名单配置),读操作完全无锁。
import java.util.concurrent.*;
CopyOnWriteArrayList<String> whitelist = new CopyOnWriteArrayList<String>();
// 写操作:底层数组整体复制,然后替换引用(代价高)
whitelist.add("192.168.1.1");
// 读 / 迭代:完全无锁,读的是调用时刻的快照,并发写入不影响当前迭代
for (String ip : whitelist) {
if (isValidIp(ip)) { /* 处理 */ }
}
| 特性 | 说明 |
|---|---|
| 适用场景 | 写操作极少(每次写要复制整个数组) |
| 读取语义 | 快照语义,迭代时不会看到并发新增的元素 |
| 迭代器 | 只读,不支持 iterator.remove() |
原子类
场景:单个变量的无锁原子操作(CAS 实现)。比 synchronized 性能更好,适合高并发计数、ID 生成、状态标志原子更新。
import java.util.concurrent.atomic.*;
// ===== 基本类型原子类 =====
AtomicInteger count = new AtomicInteger(0);
count.incrementAndGet(); // 原子 ++i,返回新值
count.getAndIncrement(); // 原子 i++,返回旧值
count.addAndGet(5); // 原子 += 5,返回新值
count.compareAndSet(10, 0); // CAS:当前值 == 10 才改为 0,返回是否成功
// ===== 高并发累加器(Java 8+,比 AtomicLong 竞争更小)=====
LongAdder adder = new LongAdder();
adder.increment(); // 分桶累加,各线程操作不同桶,减少 CAS 竞争
adder.add(100);
long total = adder.sum(); // 汇总所有桶的值(注意:有并发时 sum 不精确)
// ===== 对象引用原子替换(配置热更新场景)=====
AtomicReference<String> configRef = new AtomicReference<String>("v1");
String oldConfig = configRef.get();
String newConfig = "v2";
boolean replaced = configRef.compareAndSet(oldConfig, newConfig); // CAS 原子替换
// ===== 解决 ABA 问题:带版本号的引用 =====
// ABA 问题:值从 A 改为 B 再改回 A,普通 CAS 无法感知
AtomicStampedReference<Integer> stampedRef =
new AtomicStampedReference<Integer>(100, 0); // (value=100, stamp=0)
int[] stampHolder = new int[1];
int currentVal = stampedRef.get(stampHolder); // 同时读取 value 和 stamp
int currentStamp = stampHolder[0];
// value 和 stamp 都匹配才能更新(stamp 变化说明中间被修改过)
stampedRef.compareAndSet(currentVal, currentVal + 1, currentStamp, currentStamp + 1);
| 场景 | 推荐工具 |
|---|---|
| 高并发计数,只需最终 sum | LongAdder |
| 需要精确 CAS 的计数/版本号 | AtomicLong |
| 多个字段原子更新 | synchronized 或 Lock |
| 对象引用原子替换 | AtomicReference |
| 需要感知 ABA 的引用替换 | AtomicStampedReference |
线程池
线程是重量级对象(需调用 OS 内核 API,分配数 MB 的栈空间),频繁创建/销毁代价极高。线程池复用线程,采用生产者-消费者模式:调用方提交任务到队列,池中线程循环消费。
核心参数
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;
ThreadPoolExecutor executor = new ThreadPoolExecutor(
4, // corePoolSize:核心线程数,常驻不回收
8, // maximumPoolSize:最大线程数上限
60L, TimeUnit.SECONDS, // 非核心线程空闲多久后回收
new ArrayBlockingQueue<Runnable>(200), // workQueue(必须有界,防 OOM!)
new ThreadFactory() { // 自定义线程名,便于 jstack 排查
private final AtomicInteger seq = new AtomicInteger(0);
@Override
public Thread newThread(Runnable r) {
return new Thread(r, "biz-worker-" + seq.getAndIncrement());
}
},
new ThreadPoolExecutor.CallerRunsPolicy() // 队列满且达上限时:让提交方自己执行(降级)
);
任务提交流转:
提交任务
↓
当前线程数 < corePoolSize?
→ 是:新建核心线程立即执行
→ 否:任务队列未满?
→ 是:入队等待
→ 否:当前线程数 < maximumPoolSize?
→ 是:新建临时线程执行
→ 否:执行拒绝策略
四种拒绝策略
| 策略 | 行为 | 适用场景 |
|---|---|---|
AbortPolicy(默认) | 抛出 RejectedExecutionException | 需要感知过载 |
CallerRunsPolicy | 由提交线程自己执行任务 | 降级处理,不丢任务 |
DiscardPolicy | 直接丢弃,不抛异常 | 可接受丢失(如日志异步写) |
DiscardOldestPolicy | 丢弃队列最老的任务 | 优先保证最新任务 |
线程数设置
int cpuCount = Runtime.getRuntime().availableProcessors();
// CPU 密集型(大量计算、加密、图像处理):+1 防止线程偶尔阻塞时 CPU 空闲
int cpuBoundSize = cpuCount + 1;
// IO 密集型(数据库、HTTP 调用、文件读写):大量时间在等待 IO
// 经验值:核心数 × 2,根据监控(队列积压、CPU 利用率)动态调整
int ioBoundSize = cpuCount * 2;
实战注意事项
// ===== ⚠️ 禁止使用 Executors 快捷方法 =====
// ❌ 默认使用无界 LinkedBlockingQueue,高负载时 OOM
ExecutorService badPool = Executors.newFixedThreadPool(4);
// ❌ 无上限创建线程,突发流量时 OOM
ExecutorService badPool2 = Executors.newCachedThreadPool();
// ✅ 总是显式指定有界队列
ThreadPoolExecutor goodPool = new ThreadPoolExecutor(
4, 4, 0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<Runnable>(500),
new ThreadPoolExecutor.CallerRunsPolicy());
// ===== ⚠️ submit() 的任务内部异常会被吞掉 =====
// ❌ 这里不会打印任何异常
goodPool.execute(new Runnable() {
@Override
public void run() {
throw new RuntimeException("我被吞掉了");
}
});
// ✅ 用 Future.get() 获取任务内部异常
Future<?> future = goodPool.submit(new Runnable() {
@Override
public void run() {
throw new RuntimeException("任务内部异常");
}
});
try {
future.get();
} catch (ExecutionException e) {
System.err.println("任务失败:" + e.getCause()); // 真正的异常在 getCause() 里
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
// ===== 优雅关闭 =====
goodPool.shutdown(); // 停止接受新任务,等已提交任务执行完
try {
if (!goodPool.awaitTermination(30, TimeUnit.SECONDS)) {
goodPool.shutdownNow(); // 超过 30 秒强制停止
}
} catch (InterruptedException e) {
goodPool.shutdownNow();
Thread.currentThread().interrupt();
}
快速选型指南
flowchart TD
A["需要线程安全"] --> B{"只操作单个变量?"}
B -->|"是"| C{"有 check-then-act<br/>等复合操作?"}
C -->|"否(单步 increment)"| D["原子类<br/>AtomicXxx / LongAdder"]
C -->|"是"| E{"需要超时/中断<br/>或精确条件唤醒?"}
B -->|"否(多变量/代码块)"| E
E -->|"否"| F["synchronized"]
E -->|"是"| G["ReentrantLock"]
H["共享数据结构"] --> I{"类型"}
I -->|"Map"| J["ConcurrentHashMap"]
I -->|"List(读极多)"| K["CopyOnWriteArrayList"]
I -->|"队列/生产消费"| L["ArrayBlockingQueue"]
M["协调多线程"] --> N{"需求"}
N -->|"等待 N 个任务完成"| O["CountDownLatch"]
N -->|"N 线程同步集合点"| P["CyclicBarrier"]
N -->|"限制最大并发数"| Q["Semaphore"]
style D fill:#cfc,stroke:#060
style F fill:#cfc,stroke:#060
style G fill:#ff9,stroke:#960
style J fill:#cfc,stroke:#060
style L fill:#cfc,stroke:#060
style O fill:#e8f4f8,stroke:#4a90d9
style P fill:#e8f4f8,stroke:#4a90d9
style Q fill:#e8f4f8,stroke:#4a90d9
参考资料
- Java Tutorials — Concurrency
- java.util.concurrent 包文档(Java SE 8)
- 《Java 并发编程实战》(Brian Goetz 等著)
- Java 并发编程(JMM 原理、synchronized 锁升级、AQS 源码、ThreadLocal 深度解析)
- AQS 原理深度解析(AQS CLH 队列、state 字段、Condition 源码级解析)
评论 (0)