Java 从入门到精通(十四):线程池与异步编程——ThreadPoolExecutor、CompletableFuture 与 ForkJoin

本文默认读者已掌握线程基础与 synchronized/volatile(可参考本系列第十二篇)。JDK 以 8/11 为主,第八节虚拟线程需 JDK 21+。文中源码片段是对 OpenJDK 真实实现的简化改写:保留了位运算、CAS、锁与循环的关键结构,删掉了异常分支与注释噪音,便于阅读但不建议直接拷贝编译。所有 Demo 均可 javac 运行,建议把第三节的位运算实验与第六节的聚合 Demo 亲手跑一遍。

并发能提升吞吐,但”每来一个请求就 new Thread()“是把并发用成了灾难。第十二篇解决了”多个线程如何正确共享数据”,本篇解决另一个维度的问题:线程本身作为一种昂贵资源,该如何被复用、被管控、被观测。从 ThreadPoolExecutor 的位运算源码,到 CompletableFuture 的回调链编排,再到 JDK 21 虚拟线程对这套模型的颠覆,是一条从”会用”到”懂原理”再到”知道何时不该用”的完整路径。

一、为什么要线程池:先把线程的成本算清楚

1.1 一个线程到底有多贵

JDK 21 之前,java.lang.Thread 采用的是 1:1 线程模型:一个 Java 线程直接对应一个操作系统内核线程(pthread/lightweight process),由操作系统负责调度。这带来三项实打实的成本:

一是内存成本。 每个线程都有独立的虚拟机栈,-Xss 默认通常为 1MB(预留的虚拟地址空间,按需提交物理页),此外还有内核侧的 task_struct、内核栈与线程局部存储。粗略估算:1 万个线程仅栈的虚拟地址空间就是 10GB,实际物理占用几十到几百 MB——足以打爆一台 4G 的容器。

二是创建与销毁成本。 new Thread().start() 最终走到 pthread_create,需陷入内核、分配内核数据结构、通知调度器,实测单次开销数十微秒。若请求处理只需 100 微秒,”创建线程”就吃掉了 30% 的预算。

三是上下文切换成本。 线程数超过核数后,OS 要在多个线程间来回切换:保存/恢复寄存器、内核栈指针、页表基址,还会污染 TLB 与各级 Cache。一次切换约几千个 CPU 周期,缓存失效时可达上万。更隐蔽的是切换本身不干活:CPU 时间被”搬运工”吃掉,表现为 sys 利用率飙升、吞吐不升反降。

1.2 无限制创建线程会发生什么

下面这段代码在普通笔记本上跑几十秒就会崩溃,不同环境的报错不同,但都是灾难性的:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
import java.util.concurrent.atomic.AtomicInteger;

/**
* 反例:无限制创建线程,最终 OOM。
* 运行前请先保存好工作,可能会拖慢系统。
* 常见报错:
* java.lang.OutOfMemoryError: unable to create new native thread
* java.lang.OutOfMemoryError: Java heap space
*/
public class ThreadBombDemo {
public static void main(String[] args) {
AtomicInteger counter = new AtomicInteger();
while (true) {
new Thread(() -> {
try {
// 每个线程活 60 秒,模拟"线程还没销毁,新的又来了"
Thread.sleep(60_000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "bomb-" + counter.incrementAndGet()).start();

if (counter.get() % 500 == 0) {
System.out.println("已创建线程数:" + counter.get());
}
}
}
}

三种典型死法:其一,unable to create new native thread——触及 ulimit -u 或内核 threads-max 上限;其二,Java heap space——线程栈虽不在堆里,但 Thread 对象、线程局部变量、被引用的对象图都在堆上;其三被 OOM Killer 杀掉进程。即使侥幸没崩,几千个线程抢几个核会让 CPU 陷入”抖动”:大量时间花在切换上,延迟从毫秒劣化到秒级且毛刺不可预测——最难排查的一类线上问题。

1.3 池化:一种通用的资源治理范式

线程池的本质是资源复用 + 准入管控:预先创建一批线程(预热),用完归还而非销毁(复用),超出容量时排队或拒绝(管控),并对外暴露队列长度、活跃数等指标(可观测)。

这套范式在软件工程中反复出现,横向对比如下:

池类型 管理的资源 最小/最大容量 空闲回收 超出时的行为
线程池 OS 线程 corePoolSize / maximumPoolSize keepAliveTime 入队 / 拒绝策略
数据库连接池 TCP 连接 minIdle / maxActive minEvictableIdleTime 阻塞等待 / 抛异常
对象池 重型对象(如 PB 解析器) min / max 空闲检测 新建 / 阻塞
HTTP 连接池 长连接 maxTotal / maxPerRoute 空闲驱逐 排队等待

它们的共同点是:用有限、可控的资源去服务无限且波动的请求,并把”过载”这个事实显式暴露出来(排队、超时、拒绝),而不是假装资源无限。理解这一点,后面的拒绝策略、有界队列、动态调整就都顺理成章了。

二、ThreadPoolExecutor 的七大参数与任务决策链

2.1 七个参数的语义

ThreadPoolExecutor 最完整的构造函数有七个参数,这是理解线程池的入口:

1
2
3
4
5
6
7
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler)
参数 语义 取值依据与常见坑
corePoolSize 核心线程数,即使空闲也尽量保留(除非开启 allowCoreThreadTimeOut) 常驻规模,按负载下限设置;设为 0 时任务先进队列
maximumPoolSize 最大线程数,队列满后才会创建超出核心数的线程 只有队列有界时才有意义;无界队列下永不生效
keepAliveTime + unit 非核心线程空闲多久后被回收 潮汐业务设 30~60s;对延迟敏感业务可设更长以减少创建开销
workQueue 保存待执行任务的阻塞队列 强烈建议有界;队列容量决定了系统的缓冲能力与内存上限
threadFactory 线程创建工厂 必须自定义命名,否则 pool-1-thread-3 无法定位问题
handler 拒绝策略 生产环境建议自定义:记录日志 + 落库 + 降级

2.2 任务提交的完整决策链

提交一个任务时,线程池按以下顺序决策,这是全篇最需要记住的图(用文字描述,请默画一遍):

  1. 当前线程数 < corePoolSize → 直接创建新的核心线程执行任务(哪怕其它核心线程正空闲着)。
  2. 线程数已达 corePoolSize → 尝试把任务放进 workQueue。入队成功后还要做一次 recheck(见下文)。
  3. 队列已满 → 尝试创建”非核心线程”(直到 maximumPoolSize)来执行任务。
  4. 线程数已达 maximumPoolSize 且队列已满 → 执行拒绝策略。

对应的简化源码如下(保留关键分支):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
// ThreadPoolExecutor.execute 的三步走(简化版,保留关键分支)
public void execute(Runnable command) {
if (command == null) throw new NullPointerException();
int c = ctl.get();

// 第一步:工作线程数 < corePoolSize,直接新建核心线程,把 command 作为 firstTask
if (workerCountOf(c) < corePoolSize) {
if (addWorker(command, true))
return;
c = ctl.get(); // addWorker 失败(并发或状态变化),重新读取
}

// 第二步:线程池还在 RUNNING,尝试入队
if (isRunning(c) && workQueue.offer(command)) {
int recheck = ctl.get();
// 双重检查:入队期间可能发生了 shutdown,或所有工作线程刚好都退出了
if (!isRunning(recheck) && remove(command))
reject(command); // 已关闭 → 移除并拒绝
else if (workerCountOf(recheck) == 0)
addWorker(null, false); // 没有存活线程 → 补一个空任务线程来消费队列
}
// 第三步:入队失败(队列有界且已满)→ 尝试创建非核心线程
else if (!addWorker(command, false))
reject(command); // 达到 maximumPoolSize → 拒绝
}

2.3 为什么”先入队再加线程”是反直觉的

很多人第一次看这段代码都会疑惑:为什么满了不是先扩容线程,而是先塞进队列?队列都堆了几万个任务了,线程数还守着 core 不动,这不是”眼看要爆了还不加人”吗?

这背后是 Doug Lea 的明确设计意图,源码注释写得很清楚:线程池希望线程数保持在 corePoolSize 附近,只有当队列无法吸收任务(证明核心线程的处理能力确实不够)时,才值得付出创建线程的代价。理由有三:

  • 创建线程昂贵,而入队是相对廉价的内存操作。能用内存缓冲解决的,不要动用操作系统资源。
  • 队列是有界缓冲区,代表”系统还能承受多少积压”。队列能装下,说明只是瞬时波峰,核心线程很快就能消化。
  • 线程数增长有代价:更多线程意味着更多上下文切换与更多内存,盲目扩容反而降低吞吐。

必须记住的推论:如果队列是无界的(比如默认构造的 LinkedBlockingQueue),第 3 步永远走不到,maximumPoolSize 形同虚设。这是第五节 Executors.newFixedThreadPool 最大的坑。

顺带解释第二步的 recheck:offer 之后存在时间窗口,并发调用 shutdown() 会让线程池不再是 RUNNING,任务虽已入队但不应执行,故 remove(command) 并拒绝;另一种情况是线程数恰好掉到 0,队列里躺着任务却无人消费,须补一个 firstTask 为 null 的 worker 来”唤醒”消费。

2.4 corePoolSize 与 maximumPoolSize 怎么定

先要区分任务类型:

CPU 密集型(加密解密、图像处理、复杂计算):线程几乎不等 IO,超过核数只会增加切换开销,经验值 N + 1(N = Runtime.getRuntime().availableProcessors())。那个 +1 是为了在偶发页错误、调度抖动时有一个线程顶上,保证 CPU 不空转。

IO 密集型(数据库查询、RPC 调用、读写文件):线程大量时间在等待,经验值 2N,更严谨的公式是:

1
线程数 = N × (1 + 等待时间 / 计算时间)

比如一个请求 CPU 计算 10ms、等待下游 90ms,则线程数 = N × (1 + 90/10) = 10N。实践还可引入 Little’s Law(排队论):

1
并发数 = 目标吞吐(QPS) × 平均响应时间(RT)

若目标 1000 QPS、平均 RT 50ms,则并发度至少要 1000 × 0.05 = 50。

这些公式只是起点,不是答案。它们忽略四件事:一是锁竞争,线程再多也得串行等锁;二是容器 CPU 配额,availableProcessors() 在 K8s 里读到的是宿主机核数而非 limit;三是下游瓶颈,线程数加到 200 只是把压力传导给数据库;四是 GC。最终参数必须以压测为准:固定其它条件,单调调整线程数,观察吞吐与 P99 的拐点。

还要注意:corePoolSize == maximumPoolSize 时线程池是固定大小的;想让它”能屈能伸”,两者要拉开差距并配合有界队列使用。

三、源码剖析:ctl、execute、Worker 与回收

3.1 ctl:一个 int 装下两个状态

线程池需要同时维护运行状态与工作线程数,且两者必须原子地一起变化(否则会出现”状态改 SHUTDOWN 的同时线程数加了 1”这类竞态)。用两个 AtomicInteger 就要加锁或引入组合原子性难题,Doug Lea 的解法是:用一个 32 位 int 同时编码两个字段——高 3 位存状态,低 29 位存线程数。

为什么是 3 位?5 种状态至少需 3 位(2^3 = 8 ≥ 5),剩下 29 位存线程数,上限 (2^29) - 1 ≈ 5.37 亿,远超现实需求。

五个状态按数值单调递增排列,方便用大小比较判断”是否该停止接活”:

状态 值 含义 是否接收新任务 是否处理队列存量任务
RUNNING -1 << 29 正常运行 是 是
SHUTDOWN 0 << 29 调用 shutdown(),平和关闭 否 是
STOP 1 << 29 调用 shutdownNow(),立即停止 否 否(队列任务被丢弃并返回)
TIDYING 2 << 29 所有任务已结束,工作线程数为 0,即将执行 terminated() 否 否
TERMINATED 3 << 29 terminated() 执行完毕,彻底终结 否 否

注意 RUNNING 是负数,这样”状态是否 ≥ SHUTDOWN”这样的判断只需一次比较。相关位运算如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// ctl 的位运算:高 3 位 runState,低 29 位 workerCount
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
private static final int COUNT_BITS = Integer.SIZE - 3; // 29
private static final int CAPACITY = (1 << COUNT_BITS) - 1; // 00011111...111(低29位全1)

private static final int RUNNING = -1 << COUNT_BITS; // 111 开头
private static final int SHUTDOWN = 0 << COUNT_BITS; // 000
private static final int STOP = 1 << COUNT_BITS; // 001
private static final int TIDYING = 2 << COUNT_BITS; // 010
private static final int TERMINATED = 3 << COUNT_BITS; // 011

// 取高 3 位:用 CAPACITY 的反码做掩码
private static int runStateOf(int c) { return c & ~CAPACITY; }
// 取低 29 位
private static int workerCountOf(int c) { return c & CAPACITY; }
// 合成:状态在高位、数量在低位,直接按位或
private static int ctlOf(int rs, int wc) { return rs | wc; }

// 状态比较:c 是否处于 s 之后(数值更大 = 更接近终结)
private static boolean runStateAtLeast(int c, int s) { return c >= s; }
private static boolean isRunning(int c) { return c < SHUTDOWN; } // 只有 RUNNING 是负数

可以写个小实验验证位运算的正确性:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// 验证 ctl 编解码:动手跑一遍比看十遍注释管用
public class CtlBitsDemo {
private static final int COUNT_BITS = Integer.SIZE - 3; // 29
private static final int CAPACITY = (1 << COUNT_BITS) - 1;
private static final int RUNNING = -1 << COUNT_BITS;
private static final int SHUTDOWN = 0 << COUNT_BITS;
private static final int STOP = 1 << COUNT_BITS;

static int runStateOf(int c) { return c & ~CAPACITY; }
static int workerCountOf(int c) { return c & CAPACITY; }
static int ctlOf(int rs, int wc){ return rs | wc; }

public static void main(String[] args) {
int c = ctlOf(RUNNING, 3);
System.out.println("ctl = " + c + ",二进制 = " + Integer.toBinaryString(c));
System.out.println("线程数 = " + workerCountOf(c)); // 3
System.out.println("是否 RUNNING = " + (runStateOf(c) == RUNNING)); // true
System.out.println("容量上限 = " + CAPACITY); // 536870911

// 线程数 +1:直接对 ctl 做整数加法,低位进位不会污染高位
int c2 = c + 1;
System.out.println("+1 后线程数 = " + workerCountOf(c2) + ",状态不变 = " + (runStateOf(c2) == RUNNING));

// 状态迁移:改成 SHUTDOWN,线程数保持不变
int c3 = ctlOf(SHUTDOWN, workerCountOf(c));
System.out.println("SHUTDOWN 后线程数 = " + workerCountOf(c3) + ",isRunning = " + (c3 < SHUTDOWN));
}
}

一个精妙之处:workerCount 加 1 时直接对 ctl 做整数加法,因为低位进位最多进到第 30 位,而状态位在更高位,正常范围不会溢出污染。这使得”状态 + 数量”可一次 CAS 完成。

3.2 addWorker:CAS 自旋 + 双重检查

addWorker(firstTask, core) 负责真正创建并启动线程,分两段:先用 CAS 自旋增加 workerCount,再实例化 Worker 并启动。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
// addWorker 简化版:CAS 增加计数 + 创建 Worker 并启动
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// 状态校验:SHUTDOWN 之后不再接收新任务(但 SHUTDOWN 时允许 firstTask==null 去消费队列)
if (rs >= SHUTDOWN && !(rs == SHUTDOWN && firstTask == null && !workQueue.isEmpty()))
return false;

for (;;) {
int wc = workerCountOf(c);
// 超过容量上限或超过本次的目标边界(core ? corePoolSize : maximumPoolSize)则失败
if (wc >= CAPACITY || wc >= (core ? corePoolSize : maximumPoolSize))
return false;
// CAS 自增 workerCount,成功则跳出外层循环
if (compareAndIncrementWorkerCount(c))
break retry;
c = ctl.get(); // CAS 失败,说明有并发,重新读
if (runStateOf(c) != rs) // 状态变了,回到外层重新做状态校验
continue retry;
}
}

boolean workerStarted = false;
boolean workerAdded = false;
Worker w = new Worker(firstTask); // 封装 Thread + firstTask + AQS 锁
Thread t = w.thread;
final ReentrantLock mainLock = this.mainLock;
mainLock.lock(); // workers 是 HashSet,非线程安全,需要 mainLock 保护
try {
if (runStateOf(ctl.get()) < SHUTDOWN || (runStateOf(ctl.get()) == SHUTDOWN && firstTask == null)) {
workers.add(w);
workerAdded = true;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
t.start(); // 真正启动线程,进入 Worker.run() → runWorker()
workerStarted = true;
}
if (!workerStarted)
addWorkerFailed(w); // 回滚:从 workers 移除并 CAS 减少计数
return workerStarted;
}

要点:CAS 循环里那句 if (runStateOf(c) != rs) continue retry; 是典型的双重检查——并发下状态可能被别的线程改了,必须回外层重新校验,否则可能在 SHUTDOWN 之后错误地加入 worker。

3.3 Worker:为什么继承 AQS

Worker 是 ThreadPoolExecutor 的内部类,它同时实现了 Runnable 并继承了 AbstractQueuedSynchronizer:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
// Worker 简化版:包装线程 + 首任务 + 已完成计数;继承 AQS 实现"不可重入独占锁"
private final class Worker extends AbstractQueuedSynchronizer implements Runnable {
final Thread thread; // 真正干活的线程(由 threadFactory 创建,target 为 this)
Runnable firstTask; // 创建时顺带执行的第一个任务,可为 null
volatile long completedTasks; // 该线程完成的任务数,用于统计

Worker(Runnable firstTask) {
setState(-1); // 初始 state = -1,禁止在 runWorker 之前被中断
this.firstTask = firstTask;
this.thread = getThreadFactory().newThread(this);
}

public void run() { runWorker(this); } // 线程启动后进入主循环

// --- AQS:实现不可重入的独占锁,state 0=未锁定 1=已锁定 ---
protected boolean isHeldExclusively() { return getState() != 0; }

protected boolean tryAcquire(int unused) {
// CAS 0→1:不可重入,已在锁定状态再次 acquire 必然失败
if (compareAndSetState(0, 1)) {
setExclusiveOwnerThread(Thread.currentThread());
return true;
}
return false;
}

protected boolean tryRelease(int unused) {
setExclusiveOwnerThread(null);
setState(0);
return true;
}

public void lock() { acquire(1); }
public boolean tryLock() { return tryAcquire(1); }
public void unlock() { release(1); }
public boolean isLocked() { return isHeldExclusively(); }

// 中断"正在执行任务"的 worker 之前,先把 state 从 -1 归零,允许中断生效
void interruptIfStarted() {
Thread t;
if (getState() >= 0 && (t = thread) != null && !t.isInterrupted()) {
try { t.interrupt(); } catch (SecurityException ignore) {}
}
}
}

两个设计细节值得反复体会:

其一,为什么不用 ReentrantLock 而要自己写个不可重入的锁? ReentrantLock 可重入,shutdownNow() 里用 tryLock() 判断”线程是否空闲”时,若线程自己已持有锁(如 beforeExecute 里用户调用了 shutdownNow),tryLock() 会错误返回 true,导致正在执行的任务被中断。不可重入锁保证 tryLock() 成功 ⟺ 该 worker 确实没在执行任务。

其二,setState(-1) 与”运行中不响应中断”。 初始 state 为 -1 表示”线程尚未开始工作”,interruptIfStarted() 检查 getState() >= 0 才允许中断——这正是”Shutdown 时只有空闲线程会被中断”的实现:runWorker 一开始就 w.lock() 把 state 置为 1,shutdown 时 tryLock() 无法抢占,正在跑的任务得以执行完毕。中断只是协作式通知,不是强制杀死,这个设计保证了关闭的优雅。

3.4 runWorker 与 getTask:主循环与回收

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
// runWorker 简化版:循环从队列取任务执行,前后各留一个扩展钩子
final void runWorker(Worker w) {
Thread wt = Thread.currentThread();
Runnable task = w.firstTask;
w.firstTask = null;
w.unlock(); // state 从 -1 归 0,允许被中断(此后才进入"可被中断"状态)
boolean completedAbruptly = true;
try {
// 先执行 firstTask,之后不断 getTask();getTask 返回 null 则线程退出、被回收
while (task != null || (task = getTask()) != null) {
w.lock(); // 加锁标记"忙",shutdown 时 tryLock 失败 → 不会被中断
try {
// 如果线程池已到 STOP,确保当前线程是中断状态,让任务能感知
if ((runStateAtLeast(ctl.get(), STOP) || (Thread.interrupted() && runStateAtLeast(ctl.get(), STOP)))
&& !wt.isInterrupted())
wt.interrupt();
beforeExecute(wt, task); // 扩展点:可在子类中埋监控
try {
task.run(); // 注意:直接调用 run(),不 start 新线程
afterExecute(task, null);
} catch (Throwable ex) {
afterExecute(task, ex);
throw ex; // 抛出 → 线程死亡 → processWorkerExit 补位
}
} finally {
task = null;
w.completedTasks++;
w.unlock();
}
}
completedAbruptly = false;
} finally {
// 正常退出或异常退出都会走到这里:从 workers 移除、统计任务数、可能需要补一个线程
processWorkerExit(w, completedAbruptly);
}
}

重点提醒:任务抛出的异常会导致当前工作线程直接死亡,再由 processWorkerExit 补充新 worker。线程池”吞掉”了任务异常(用 execute 时)但不会崩溃,代价是线程重建。所以建议任务里用 try-catch 包住业务逻辑,并重写 afterExecute 记录日志。

getTask() 决定了线程何时被回收:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
// getTask 简化版:决定线程是阻塞等待、超时等待,还是直接返回 null 被回收
private Runnable getTask() {
boolean timedOut = false;
for (;;) {
int c = ctl.get();
// 线程池已停止时,把 workerCount 减到 0 并返回 null,让线程退出
if (runStateAtLeast(c, SHUTDOWN) && (runStateAtLeast(c, STOP) || workQueue.isEmpty())) {
decrementWorkerCount();
return null;
}
int wc = workerCountOf(c);
// 关键:是否允许核心线程超时。allowCoreThreadTimeOut=true 时,核心线程也会超时退出
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;

if ((wc > maximumPoolSize || (timed && timedOut)) && (wc > 1 || workQueue.isEmpty())) {
if (compareAndDecrementWorkerCount(c)) // CAS 减少计数后返回 null → 线程消亡
return null;
continue;
}
try {
// 超时等待 poll(keepAlive) vs 无限阻塞 take()
Runnable r = timed ? workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS)
: workQueue.take();
if (r != null) return r;
timedOut = true; // 本轮超时,下一轮循环会判定回收
} catch (InterruptedException retry) {
timedOut = false; // 被中断(多半是 shutdown)→ 重新判断状态
}
}
}

由此得到回收规则:只有”非核心线程”(wc > corePoolSize)或开启了 allowCoreThreadTimeOut 的线程,才会在 keepAliveTime 空闲后被回收;否则线程一直阻塞在 workQueue.take() 上。allowCoreThreadTimeOut(true) 的价值是让线程池在潮汐流量下真正缩容,避免低峰期白占线程与栈内存。

3.5 shutdown 与 shutdownNow:一字之差,行为天壤

对比项 shutdown() shutdownNow()
状态迁移 RUNNING/SHUTDOWN → SHUTDOWN → STOP
是否接收新任务 否 否
队列里积压的任务 继续执行完毕 丢弃,以 List<Runnable> 返回
正在执行的任务 允许跑完(不中断) 尝试 interrupt()(协作式,取决于任务是否响应中断)
中断范围 只中断空闲 worker(tryLock 成功的) 中断所有已启动的 worker
返回值 void List<Runnable> 未执行任务
适用场景 优雅停机:注册中心下线后等存量请求处理完 紧急止损
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
import java.util.List;
import java.util.concurrent.*;

/** shutdown 与 shutdownNow 的行为差异 */
public class ShutdownDemo {
public static void main(String[] args) throws Exception {
ThreadPoolExecutor pool = new ThreadPoolExecutor(
1, 1, 0, TimeUnit.SECONDS, new LinkedBlockingQueue<>(10),
r -> new Thread(r, "worker-1"));

// 提交一个 3 秒的长任务 + 3 个队列任务
pool.execute(() -> sleep(3000));
for (int i = 0; i < 3; i++) {
final int n = i;
pool.execute(() -> System.out.println("执行队列任务 " + n));
}
Thread.sleep(100);

System.out.println("--- shutdownNow:立即停止 ---");
List<Runnable> dropped = pool.shutdownNow(); // 返回被丢弃的队列任务
System.out.println("被丢弃未执行的任务数 = " + dropped.size()); // 3
System.out.println("awaitTermination = " + pool.awaitTermination(5, TimeUnit.SECONDS));

// 对比:若这里换成 shutdown(),则 3 个队列任务都会被正常执行完
}

static void sleep(long ms) {
try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
}

优雅停机的标准写法是三段式:shutdown() → awaitTermination(timeout) → 未结束则 shutdownNow() 并二次等待。Spring Boot 里可给线程池 bean 设置 setWaitForTasksToCompleteOnShutdown(true) 与 setAwaitTerminationSeconds(30)。

3.6 状态迁移全景

五个状态的迁移路径只有四条,且不可逆:

  • RUNNING → SHUTDOWN:调用 shutdown()。
  • RUNNING/SHUTDOWN → STOP:调用 shutdownNow()。
  • SHUTDOWN → TIDYING:队列已空且工作线程数为 0。
  • STOP → TIDYING:工作线程数为 0。
  • TIDYING → TERMINATED:terminated() 钩子执行完毕。

tryTerminate() 是唯一的推进器:每当 worker 退出、任务被移除都会尝试调用它;发现”SHUTDOWN 且队列空且 wc==0”(或 STOP 且 wc==0)时,就把状态 CAS 推进到 TIDYING,执行 terminated(),再置为 TERMINATED,最后 termination.signalAll() 唤醒阻塞在 awaitTermination 上的线程。注意:terminated() 是空实现,正是留给我们的资源清理扩展点。

四、工作队列与拒绝策略

4.1 四种队列如何改变线程池的性格

队列不是被动容器,它直接决定线程池的行为曲线:

队列 结构 是否有界 对线程池的影响
ArrayBlockingQueue 数组 + 单把 ReentrantLock 有界(构造时必填容量) 队列满即触发扩容到 max,行为可预测;单锁下生产消费互斥,吞吐略低
LinkedBlockingQueue 链表 + 入队出队两把锁 默认无界(Integer.MAX_VALUE),可显式指定容量 无界时 maximumPoolSize 失效,堆积可致 OOM;有界时吞吐优于 ArrayBlockingQueue
SynchronousQueue 不存储元素,生产者必须等到消费者接手 容量恒为 0 任务直接移交,无缓冲 → 必须立刻开新线程,是 newCachedThreadPool 高并发吞吐的来源
PriorityBlockingQueue 堆结构,按 Comparator 排序 无界 支持优先级调度,但无界同样有 OOM 风险;优先级反转需自行处理
DelayedWorkQueue 堆 + 延迟 无界 ScheduledThreadPoolExecutor 专用,按到期时间出队

关于 LinkedBlockingQueue 的陷阱值得单独强调:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// 同一个 core/max 配置,仅仅换一个队列,行为天差地别
public class QueueImpactDemo {
public static void main(String[] args) {
// A:无界队列 —— maximumPoolSize 永远是摆设,第 3 步永远走不到
ThreadPoolExecutor unbounded = new ThreadPoolExecutor(
2, 16, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>()); // 容量 = Integer.MAX_VALUE

// B:有界队列 —— 队列满 100 之后才会创建第 3~16 个线程,再满则拒绝
ThreadPoolExecutor bounded = new ThreadPoolExecutor(
2, 16, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100)); // 显式容量
}
}

Executors.newFixedThreadPool(n) 内部用的正是 A 方案,这就是为什么它在流量洪峰下会静默堆积、最终 Java heap space。结论:生产环境永远显式指定队列容量。

4.2 四种内置拒绝策略

策略 行为 优点 风险 适用
AbortPolicy(默认) 抛 RejectedExecutionException 快速失败,问题立刻暴露 调用方不捕获会导致请求链路中断 需要强告警的核心链路
CallerRunsPolicy 由提交任务的线程自己执行 天然反压,削峰填谷 占用业务线程,可能拖垮上游;线程池已关闭时静默丢弃 可接受降级的批处理、后台任务
DiscardPolicy 静默丢弃,无任何提示 实现简单 数据悄悄丢失,最难排查 可丢弃的埋点/日志上报
DiscardOldestPolicy 丢弃队列队头(最老)任务,再尝试提交 优先处理新任务 可能丢弃关键的老任务 实时性优先的场景

CallerRunsPolicy 的反压效果值得理解:线程池饱和时提交者被迫自己干活,于是没时间再提交新任务,提交速率自然下降——这是一个负反馈闭环,让系统在过载时”减速”而非”崩溃”。但它有两个副作用:一是若提交者是 Tomcat 的 HTTP 工作线程或 Netty 的 EventLoop,反压会直接传导为接口超时或事件循环阻塞;二是 JDK 8 实现中若线程池已 shutdown,策略会静默丢弃任务而不抛异常。

4.3 自定义拒绝策略:日志 + 落库 + 降级

生产环境几乎不应该直接用内置策略,至少要能”留痕”:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;

/**
* 生产级拒绝策略:记日志(含线程池运行状态)+ 落库待补偿 + 触发告警
*/
public class AlertingRejectHandler implements RejectedExecutionHandler {
private final String poolName;
private final AtomicLong rejectedCount = new AtomicLong();
// 实际项目中注入 DAO / MQ / 告警客户端
// private final TaskRepository repository;

public AlertingRejectHandler(String poolName) { this.poolName = poolName; }

@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
long n = rejectedCount.incrementAndGet();
// 1)日志:带上关键指标,便于事后还原现场
System.err.printf("[REJECT] pool=%s seq=%d active=%d poolSize=%d queue=%d/%d completed=%d%n",
poolName, n, e.getActiveCount(), e.getPoolSize(),
e.getQueue().size(), e.getQueue().remainingCapacity() + e.getQueue().size(),
e.getCompletedTaskCount());

// 2)落库 / 投递到 MQ,等待补偿任务重放(保证任务不丢)
// repository.saveFailedTask(TaskStatus.REJECTED, r.toString());

// 3)降级:尝试由调用线程执行(等效 CallerRunsPolicy),但要防止拖垮上游
// 这里做个判断:只在线程池仍在运行时才让调用方兜底
if (!e.isShutdown()) {
r.run();
}

// 4)告警:按 n % 100 == 1 之类的频率做抑制,避免告警风暴
if (n % 100 == 1) {
System.err.println("[ALERT] 线程池 " + poolName + " 已拒绝 " + n + " 个任务,请扩容或限流!");
}
}
}

真正稳妥的做法是不要依赖拒绝策略兜底:在入口用 Sentinel / Resilience4j 做限流,把过载挡在外面,拒绝策略只作最后一道保险。

五、Executors 的陷阱与正确姿势

5.1 四个工厂方法的真实参数

工厂方法 core max 队列 核心风险
newFixedThreadPool(n) n n LinkedBlockingQueue 无界 任务堆积 → OutOfMemoryError: Java heap space
newSingleThreadExecutor() 1 1 LinkedBlockingQueue 无界 同上;且被 FinalizableDelegatedExecutorService 包装,无法强转为 ThreadPoolExecutor 调参
newCachedThreadPool() 0 Integer.MAX_VALUE SynchronousQueue 线程数无上限 → 高并发下线程爆炸、unable to create new native thread
newScheduledThreadPool(n) n Integer.MAX_VALUE DelayedWorkQueue 延时任务堆积时同样可能爆线程;且定时任务异常会静默终止后续调度
newWorkStealingPool() — — 基于 ForkJoinPool 默认并行度为 CPU 核数;IO 密集任务会阻塞 worker

这就是《阿里巴巴 Java 开发手册》强制要求”线程池不允许使用 Executors 创建”的原因——Executors 把最危险的两个参数(无界队列、无限最大线程数)藏在便利性的外衣下,让资源耗尽的风险从编译期推迟到线上爆炸时。

5.2 手动创建的标准模板

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/** 生产环境线程池创建标准模板 */
public class PoolFactory {
public static ThreadPoolExecutor create(String bizName, int core, int max,
int queueCapacity, long keepAliveSec) {
// 1)自定义 ThreadFactory:命名 + 非守护 + 异常处理器 + 明确栈大小
ThreadFactory factory = new ThreadFactory() {
private final AtomicInteger idx = new AtomicInteger(1);
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, bizName + "-pool-" + idx.getAndIncrement());
t.setDaemon(false); // 业务池必须非守护,否则 JVM 退出会丢任务
t.setUncaughtExceptionHandler((th, e) ->
System.err.println("线程 " + th.getName() + " 发生未捕获异常:" + e.getMessage()));
return t;
}
};

// 2)有界队列 + 3)自定义拒绝策略 + 4)明确的 keepAlive
return new ThreadPoolExecutor(
core, max, keepAliveSec, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(queueCapacity),
factory,
new AlertingRejectHandler(bizName));
}

// 拆分线程池是"舱壁隔离"的关键:慢接口不能拖垮快接口
public static void main(String[] args) {
ThreadPoolExecutor orderPool = create("order", 8, 16, 500, 60);
ThreadPoolExecutor notifyPool = create("notify", 2, 4, 200, 30);
orderPool.execute(() -> System.out.println("订单任务执行于 " + Thread.currentThread().getName()));
notifyPool.execute(() -> System.out.println("通知任务执行于 " + Thread.currentThread().getName()));
orderPool.shutdown();
notifyPool.shutdown();
}
}

四条铁律:必须命名线程(pool-3-thread-8 在 jstack 里毫无信息量);必须有界队列(宁可拒绝,不可 OOM);必须自定义拒绝策略(至少留痕);必须按业务隔离线程池(慢 SQL 拖垮共用池,所有业务陪葬)。

5.3 动态调参与监控

线程池参数不必写死,可在运行时调整,配合配置中心(Nacos / Apollo)实现”高峰扩容、低峰缩容”:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
import java.util.concurrent.*;

/** 线程池动态调参 + 指标暴露 */
public class DynamicPoolTuning {
public static void main(String[] args) {
ThreadPoolExecutor pool = new ThreadPoolExecutor(
4, 8, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(1000),
r -> new Thread(r, "dyn-pool"));

// 运行时调参:setCorePoolSize / setMaximumPoolSize 均线程安全,立即生效
pool.setCorePoolSize(16);
pool.setMaximumPoolSize(32);
pool.setKeepAliveTime(120, TimeUnit.SECONDS);
// 开启后核心线程也会超时回收,低峰期可缩容到 0(注意 core 不能为 0 且要保证有线程消费队列)
pool.allowCoreThreadTimeOut(true);

// 注意:setCorePoolSize(新值 < 旧值) 时,多余的核心线程会在下次 getTask 超时后被回收;
// setCorePoolSize(新值 > 旧值) 时,若队列非空,会立即创建新线程去消费积压任务。

// 监控指标:建议每 10~60 秒采集一次并上报 Micrometer / Prometheus
printMetrics(pool);
pool.shutdown();
}

static void printMetrics(ThreadPoolExecutor pool) {
System.out.printf("poolSize=%d(core=%d,max=%d) active=%d queue=%d completed=%d taskTotal=%d%n",
pool.getPoolSize(), pool.getCorePoolSize(), pool.getMaximumPoolSize(),
pool.getActiveCount(), pool.getQueue().size(),
pool.getCompletedTaskCount(), pool.getTaskCount());
}
}

更精细的监控可继承 ThreadPoolExecutor 重写三个钩子:beforeExecute(记录开始时间到 ThreadLocal)、afterExecute(计算耗时与异常)、terminated(关闭时打印汇总)。把 P50/P99 耗时、拒绝次数、队列水位打点到 Micrometer,并配置”队列使用率 > 80% 持续 1 分钟”告警,就能在雪崩前发现问题。

六、CompletableFuture:把异步从”取结果”升级为”编排流程”

6.1 Future 的三个痛点

FutureTask 解决了”异步执行 + 稍后取结果”,但用它编排多个异步任务极其痛苦(四个痛点见代码注释):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
import java.util.concurrent.*;

/** Future 的局限:get 阻塞、无法组合、异常被包装 */
public class FutureLimitDemo {
public static void main(String[] args) throws Exception {
ExecutorService pool = Executors.newFixedThreadPool(2);

Future<String> f1 = pool.submit(() -> queryUser());
Future<String> f2 = pool.submit(() -> queryOrder());

// 痛点 1:get() 阻塞当前线程 —— 想"两个都完成再做"只能顺序阻塞,无法声明式组合
String u = f1.get();
String o = f2.get(); // 即便 f2 先完成,也要等 f1

// 痛点 2:无法把 f1 的结果直接传给下一个异步任务,只能自己写回调或再 submit
// 痛点 3:异常被包装成 ExecutionException,需要层层 unwrap
// 痛点 4:没有超时 + 默认值的组合能力(get(timeout) 超时抛异常,还得自己写降级)

System.out.println(u + " | " + o);
pool.shutdown();
}

static String queryUser() { sleep(100); return "user:张三"; }
static String queryOrder() { sleep(150); return "order:A001"; }
static void sleep(long ms) { try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }
}

CompletableFuture 正是为把”异步”从取值升级为编排而来:它既能被手动完成(complete),又能挂接一串回调,还能对多个未来做组合与异常处理。

6.2 核心模型与创建方式

内部只有两个字段撑起全部能力:volatile Object result(结果或包装的异常)与 volatile Completion stack(一个 Treiber 栈,保存所有依赖它的回调)。result 写入后弹出栈中所有回调并逐个触发——典型的”完成即通知”。

分类 方法 说明
创建 supplyAsync(Supplier<U>) 有返回值,默认用 ForkJoinPool.commonPool
supplyAsync(sup, executor) 推荐:显式指定业务线程池
runAsync(Runnable[, executor]) 无返回值
completedFuture(value) 包装一个已完成的结果
串行 thenApply(Function) 同步转换结果 T → U,返回新 CF
thenApplyAsync(fn[, ex]) 异步转换,切换到指定线程池执行
thenAccept(Consumer) 消费结果,返回 CF<Void>
thenRun(Runnable) 不关心结果,只在其后执行
thenCompose(Function) 扁平化:避免 CF<CF<T>>,等价于”异步版 thenApply”
组合 thenCombine(other, BiFunction) 两个都完成后合并结果
applyToEither(other, Function) 谁先完成用谁的结果
allOf(cfs...) 全部完成(返回 CF<Void>,需自己取值)
anyOf(cfs...) 任一完成
异常 exceptionally(Function) 只在异常时给替代值,正常时透传
handle(BiFunction) 结果/异常都能拿到,可改变返回值
whenComplete(BiConsumer) 类似 finally,不改变结果,异常继续传播
超时 orTimeout(t, unit) JDK 9+,超时抛 CompletionException(TimeoutException)
completeOnTimeout(v, t, unit) JDK 9+,超时填充默认值

为什么必须指定线程池? supplyAsync 不传 executor 时用 ForkJoinPool.commonPool,其并行度默认 CPU 核数 - 1,且 JVM 全局共享。后果有二:一是 IO 密集任务长时间占用 commonPool 线程,导致 parallelStream 与其它 CompletableFuture 全部饿死;二是 commonPool 线程是守护线程,主线程退出后异步任务被直接放弃。

6.3 串行:thenApply / thenCompose / thenAccept 的差异

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
import java.util.concurrent.*;

/** 串行编排:thenApply / thenCompose / thenAccept / 异步回调线程归属 */
public class SerialChainDemo {
public static void main(String[] args) throws Exception {
ExecutorService pool = Executors.newFixedThreadPool(4, r -> new Thread(r, "biz-" + r.hashCode()));

// 1)thenApply:同步转换,由"上一个任务完成的线程"继续执行
CompletableFuture<String> f = CompletableFuture
.supplyAsync(() -> "订单A", pool)
.thenApply(s -> s + "-已支付") // 同一个 biz 线程继续跑
.thenApply(String::toUpperCase);

// 2)thenCompose:把 "返回 CF 的函数" 扁平化,避免 CF<CF<T>>
CompletableFuture<String> flat = CompletableFuture
.supplyAsync(() -> 1001L, pool) // CF<Long>
.thenCompose(id -> queryDetailAsync(id, pool)); // Long -> CF<String>,展开为 CF<String>

// 若误用 thenApply,类型会变成 CompletableFuture<CompletableFuture<String>>
CompletableFuture<CompletableFuture<String>> nested = CompletableFuture
.supplyAsync(() -> 1001L, pool)
.thenApply(id -> queryDetailAsync(id, pool));

// 3)thenAccept / thenRun:只消费不产生新值
CompletableFuture<Void> done = flat.thenAccept(v -> System.out.println("结果:" + v))
.thenRun(() -> System.out.println("链路结束"));

System.out.println("thenApply 结果 = " + f.get());
System.out.println("thenCompose 结果 = " + flat.get());
System.out.println("嵌套类型结果 = " + nested.get().get()); // 要 get 两次,很别扭
done.join();
pool.shutdown();
}

static CompletableFuture<String> queryDetailAsync(long id, Executor ex) {
return CompletableFuture.supplyAsync(() -> "detail-" + id, ex);
}
}

判断口诀:转换逻辑是同步计算用 thenApply,转换逻辑本身要异步(返回 CompletableFuture)用 thenCompose。

6.4 异常处理:exceptionally / handle / whenComplete 怎么选

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
import java.util.concurrent.*;

/** 三种异常处理方式的差异 */
public class ExceptionHandleDemo {
public static void main(String[] args) throws Exception {
ExecutorService pool = Executors.newFixedThreadPool(2);

// exceptionally:只处理异常,正常时透传。相当于 catch → 返回默认值
CompletableFuture<String> f1 = CompletableFuture.<String>supplyAsync(() -> {
throw new RuntimeException("下游超时");
}, pool).exceptionally(ex -> "默认用户");

// handle:无论成功失败都进入,可以做统一转换,返回值会替换原结果
CompletableFuture<String> f2 = CompletableFuture.<String>supplyAsync(() -> {
throw new RuntimeException("库存服务不可用");
}, pool).handle((res, ex) -> res != null ? res : "降级结果(" + ex.getMessage() + ")");

// whenComplete:类似 finally,不改变结果;如果这里不处理,异常会继续往下游抛
CompletableFuture<String> f3 = CompletableFuture.supplyAsync(() -> "ok", pool)
.whenComplete((res, ex) -> System.out.println("清理资源,res=" + res + ", ex=" + ex));

System.out.println(f1.get()); // 默认用户
System.out.println(f2.get()); // 降级结果(...)
System.out.println(f3.get()); // ok

// 注意:未处理的异常在调用 get() 时会被包装成 ExecutionException
CompletableFuture<String> f4 = CompletableFuture.<String>supplyAsync(() -> {
throw new IllegalStateException("boom");
}, pool);
try {
f4.get();
} catch (ExecutionException e) {
System.out.println("捕获到真实原因:" + e.getCause()); // IllegalStateException: boom
}
pool.shutdown();
}
}

选择建议:需要兜底默认值用 exceptionally;需要把成功/失败统一映射成同一种返回结构(如 Result<T>)用 handle;需要记录日志、关闭资源但不想影响链路用 whenComplete。另外 join() 抛非受检的 CompletionException,在 lambda 里比 get() 顺手。

6.5 实战:并行调用三个下游接口并聚合

最常见场景:详情页要同时拉商品、库存、营销三个接口,串行 400ms,并行只要 200ms。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
import java.util.concurrent.*;
import java.util.function.Supplier;

/** 实战:三接口并行聚合 + 超时降级 + 线程池隔离 */
public class AggregateDemo {
// 关键:为三个下游分别准备线程池,实现舱壁隔离,避免一个慢服务拖垮整条链路
static final ExecutorService ITEM_POOL = newFixedPool("item", 8);
static final ExecutorService STOCK_POOL = newFixedPool("stock", 8);
static final ExecutorService PROMO_POOL = newFixedPool("promo", 8);

public static void main(String[] args) {
long start = System.currentTimeMillis();

CompletableFuture<String> itemF = callAsync(() -> rpc("商品服务", 120), "商品", ITEM_POOL);
CompletableFuture<String> stockF = callAsync(() -> rpc("库存服务", 80), "库存", STOCK_POOL);
CompletableFuture<String> promoF = callAsync(() -> rpc("营销服务", 200), "营销", PROMO_POOL);

// allOf 本身返回 CF<Void>,需要自己从各个 Future 里取值
CompletableFuture<Detail> detail = CompletableFuture.allOf(itemF, stockF, promoF)
.thenApply(v -> new Detail(itemF.join(), stockF.join(), promoF.join()));

// 整体再加一层超时兜底
Detail result = detail
.completeOnTimeout(Detail.partial(itemF.getNow("商品-未知"),
stockF.getNow("库存-未知"),
promoF.getNow("营销-无")),
300, TimeUnit.MILLISECONDS) // 超过 300ms 用已有结果兜底
.exceptionally(ex -> Detail.fallback("聚合失败:" + ex.getMessage()))
.join();

System.out.println(result);
System.out.println("总耗时 = " + (System.currentTimeMillis() - start) + "ms");
shutdownAll();
}

/** 单个调用:带超时降级,任一服务超时不影响整体 */
static <T> CompletableFuture<T> callAsync(Supplier<T> supplier, String name, ExecutorService pool) {
return CompletableFuture.supplyAsync(supplier, pool)
.orTimeout(150, TimeUnit.MILLISECONDS) // JDK 9+
.exceptionally(ex -> (T) (name + "-降级")); // 超时或异常都降级
}

static ExecutorService newFixedPool(String name, int n) {
return Executors.newFixedThreadPool(n, r -> new Thread(r, name + "-" + r.hashCode()));
}

static String rpc(String svc, int costMs) {
try { Thread.sleep(costMs); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
return svc + "数据";
}

static void shutdownAll() {
ITEM_POOL.shutdown(); STOCK_POOL.shutdown(); PROMO_POOL.shutdown();
}

record Detail(String item, String stock, String promo) {
static Detail partial(String i, String s, String p) { return new Detail(i, s, p); }
static Detail fallback(String msg) { return new Detail(msg, msg, msg); }
}
}

这段代码有四个可借鉴的点:一是线程池隔离,库存慢不会耗尽商品的线程;二是 per-call 超时,用 orTimeout 给每个下游单独设阈值;三是分层降级,单个服务降级到字符串、整体再兜一层 completeOnTimeout;四是 getNow,能在任务未完成时立刻返回默认值,是构造”部分结果”的关键。

七、ForkJoinPool 与并行流

7.1 分治与工作窃取

ForkJoinPool 为分治型(divide and conquer)计算任务而设计,核心创新是工作窃取(work-stealing):每个工作线程维护一个双端队列(WorkQueue)。

  • 自己产生的子任务通过 fork() push 到队头,自己执行时也从队头 pop(LIFO)。
  • 自己的队列空了,就随机挑一个其它线程的队列,从队尾偷(FIFO)。

“自己 LIFO、偷别人 FIFO”有两层考量:LIFO 取自己的任务能拿到最新(粒度更小、数据还在缓存里)的任务,局部性更好;从队尾窃取拿到的是最老的任务,而最老的任务往往”最大”,一次窃取就获得足够工作,减少窃取次数与竞争。同时窃取在队尾、消费在队头,多数情况下操作队列的不同端,冲突概率极低。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
import java.util.concurrent.*;

/** RecursiveTask:大数组并行求和 */
public class ForkJoinSumDemo {
static final int THRESHOLD = 10_000; // 阈值:小于它就直接算,避免过度拆分

public static void main(String[] args) {
long[] array = new long[100_000_000];
for (int i = 0; i < array.length; i++) array[i] = i + 1;

// 方式一:使用 commonPool(ForkJoinPool.commonPool())
long r1 = new ForkJoinPool().invoke(new SumTask(array, 0, array.length));
// 顺便对比一下串行
long serial = 0;
for (long v : array) serial += v;
System.out.println("ForkJoin = " + r1 + ",串行 = " + serial);
}

static class SumTask extends RecursiveTask<Long> {
private final long[] arr;
private final int start, end;

SumTask(long[] arr, int start, int end) { this.arr = arr; this.start = start; this.end = end; }

@Override
protected Long compute() {
int len = end - start;
if (len <= THRESHOLD) { // 足够小,直接计算(这是防止过度拆分的关键)
long sum = 0;
for (int i = start; i < end; i++) sum += arr[i];
return sum;
}
int mid = start + (len >>> 1);
SumTask left = new SumTask(arr, start, mid);
SumTask right = new SumTask(arr, mid, end);
left.fork(); // 左半异步执行(push 到自己的 WorkQueue)
long rightResult = right.compute(); // 右半同步执行(复用当前线程,减少一次调度)
long leftResult = left.join(); // 等待左半结果
return leftResult + rightResult;
}
}
}

注意 left.fork(); right.compute(); left.join(); 这个惯用写法——一半异步、一半同步,比”两个都 fork 再 join”少一次任务调度,是官方推荐模式。

再看归并排序,更能体现分治思想:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
import java.util.Arrays;
import java.util.concurrent.*;

/** RecursiveAction:并行归并排序(无返回值) */
public class ForkJoinSortDemo {
public static void main(String[] args) {
int[] data = new int[1_000_000];
for (int i = 0; i < data.length; i++) data[i] = (int) (Math.random() * 1_000_000);

ForkJoinPool pool = new ForkJoinPool(Runtime.getRuntime().availableProcessors());
pool.invoke(new SortTask(data, new int[data.length]));
pool.shutdown();

System.out.println("是否有序 = " + isSorted(data));
System.out.println("前 10 个 = " + Arrays.toString(Arrays.copyOf(data, 10)));
}

static class SortTask extends RecursiveAction {
private final int[] src;
private final int[] tmp;
private final int lo, hi;
private static final int THRESHOLD = 5_000;

SortTask(int[] src, int[] tmp, int lo, int hi) { this.src = src; this.tmp = tmp; this.lo = lo; this.hi = hi; }

@Override
protected void compute() {
if (hi - lo <= THRESHOLD) {
Arrays.sort(src, lo, hi); // 小规模直接交给双轴快排,比继续拆分更快
return;
}
int mid = (lo + hi) >>> 1;
SortTask left = new SortTask(src, tmp, lo, mid);
SortTask right = new SortTask(src, tmp, mid, hi);
invokeAll(left, right); // 等价 fork+fork+join+join,语义更清晰
merge(src, tmp, lo, mid, hi);
}

private void merge(int[] a, int[] tmp, int lo, int mid, int hi) {
System.arraycopy(a, lo, tmp, lo, hi - lo);
int i = lo, j = mid, k = lo;
while (i < mid && j < hi) {
a[k++] = (tmp[i] <= tmp[j]) ? tmp[i++] : tmp[j++];
}
while (i < mid) a[k++] = tmp[i++];
while (j < hi) a[k++] = tmp[j++];
}
}

static boolean isSorted(int[] a) {
for (int i = 1; i < a.length; i++) if (a[i - 1] > a[i]) return false;
return true;
}
}

7.2 parallelStream 的底层与”不该用”的四种情况

stream().parallel() 底层就是 ForkJoinPool.commonPool(),并行度默认 availableProcessors() - 1(主线程也算一份),可通过 -Djava.util.concurrent.ForkJoinPool.common.parallelism=N 调整。拆分由 Spliterator 完成:数组、ArrayList 拆分均匀效率高;LinkedList、BufferedReader.lines() 拆分极差,几乎无收益。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
import java.util.*;
import java.util.stream.*;

/** parallelStream 的正确与错误用法 */
public class ParallelStreamDemo {
public static void main(String[] args) {
List<Integer> nums = new ArrayList<>();
for (int i = 1; i <= 1_000_000; i++) nums.add(i);

// 正确:无状态、纯 CPU 计算、数据源易拆分、结果用 reduce 聚合 → 有加速
long sum = nums.parallelStream().mapToLong(Integer::longValue).sum();
System.out.println("sum = " + sum);

// 错误 1:共享可变状态。ArrayList 非线程安全 → 结果丢失甚至抛异常
List<Integer> unsafe = new ArrayList<>();
nums.parallelStream().forEach(unsafe::add); // 数据量小时"看起来没问题",上线就翻车
System.out.println("错误用法结果数 = " + unsafe.size() + ",期望 = " + nums.size());

// 正确写法:用 collect 归约,让框架处理线程安全
List<Integer> safe = nums.parallelStream().filter(x -> x % 2 == 0).toList();
System.out.println("安全结果数 = " + safe.size());

// 错误 2:用 forEach + 外部锁,把并行退化成串行,还多付了调度成本
// 错误 3:IO 密集任务(数据库查询、HTTP 调用)放并行流 → 占满 commonPool,影响全 JVM
}
}

不该用并行流的四种情况:① 数据量小(万级以下,拆分合并开销大于收益);② IO 密集(线程阻塞在等待上,反而占用全局 commonPool);③ 有状态 lambda 或依赖顺序(forEachOrdered、limit 会严重削弱并行度);④ 数据源不易拆分或装箱成本高(Stream<Integer> 拆装箱开销可观,优先 IntStream)。此外不要在并行流里做同步阻塞,也不要修改共享变量——用 collect/reduce 表达归约意图。

八、JDK 21 虚拟线程与最佳实践清单

8.1 平台线程的成本与虚拟线程模型

前面所有内容都建立在”线程很贵”这个前提上。JDK 21(JEP 444)带来的虚拟线程(Virtual Thread)正是为打破这个前提:它是由 JVM 管理、不直接对应 OS 线程的轻量级线程,创建成本在微秒级、内存占用从 MB 级降到 KB 级,因此可创建百万级而无压力。

底层模型有三个关键角色:

  • Continuation:可暂停/恢复的计算单元。虚拟线程的调用栈不再固定在 OS 栈上,而以栈帧形式保存在堆内存里,阻塞时可 yield 让出载体线程,之后再 run 恢复。
  • 载体线程(Carrier Thread):真正执行虚拟线程字节码的平台线程。默认调度器是一个 ForkJoinPool,并行度等于 CPU 核数。
  • 调度器(Scheduler):负责把虚拟线程挂载(mount)到载体线程上。

关键机制:虚拟线程执行阻塞操作(LockSupport.park、Socket 读写、Thread.sleep)时,JVM 会把它改写成一次 Continuation.yield——先卸载(unmount)虚拟线程、把栈帧搬回堆,释放载体线程去执行其它虚拟线程;条件满足后再挂载(mount)到空闲载体线程继续执行。对 Java 代码而言 Thread.sleep() 仍是阻塞的,但底层 OS 线程没有闲着——这就是”用同步的写法获得异步的性能”。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
import java.time.Duration;
import java.util.concurrent.*;
import java.util.stream.IntStream;

/** JDK 21 虚拟线程:三种创建方式 + 百万级并发演示(需 JDK 21+) */
public class VirtualThreadDemo {
public static void main(String[] args) throws Exception {
// 方式一:Thread.ofVirtual() 构建器
Thread t1 = Thread.ofVirtual().name("vt-1").start(() -> System.out.println("hello " + Thread.currentThread()));
t1.join();

// 方式二:快捷方法 startVirtualThread
Thread.startVirtualThread(() -> System.out.println("快捷创建:" + Thread.currentThread().isVirtual()));

// 方式三:Executor —— 每个任务一个虚拟线程(注意:虚拟线程不要池化!)
try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) {
long start = System.currentTimeMillis();
// 10 万个"阻塞 100 毫秒"的任务:换成平台线程池需要巨大线程数,虚拟线程轻松完成
IntStream.range(0, 100_000).forEach(i ->
executor.submit(() -> {
Thread.sleep(Duration.ofMillis(100));
return i;
}));
System.out.println("提交 10 万任务耗时 = " + (System.currentTimeMillis() - start) + "ms");
}
}
}

8.2 适用场景与四条禁忌

适用场景:高并发 IO 密集的服务端程序——典型的”每请求一线程”Web 服务。过去 200 并发就要 200 个平台线程,现在 10 万并发也只是 10 万个虚拟线程,代码依然是同步写法。Spring Boot 3.2+ 只需 spring.threads.virtual.enabled=true 即可让 Tomcat 使用虚拟线程。

禁忌一:synchronized 造成的钉扎(pinning)。 JDK 21~23 中,虚拟线程在 synchronized 块内阻塞时无法卸载——监视器持有者记录在 OS 线程层面,JVM 无法安全地把栈搬走,结果载体线程被钉住,并发度退化。可用 -Djdk.tracePinnedThreads=short 打印钉扎堆栈定位;解法是换成 ReentrantLock(JUC 的锁已适配虚拟线程)。JEP 491(JDK 24)重新实现了 synchronized,使其不再钉扎。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
import java.util.concurrent.locks.ReentrantLock;

/** 钉扎对比:synchronized 会 pin 住载体线程,ReentrantLock 不会 */
public class PinningDemo {
private static final Object MONITOR = new Object();
private static final ReentrantLock LOCK = new ReentrantLock();

public static void main(String[] args) throws Exception {
// 反例:synchronized 内阻塞 → 虚拟线程无法卸载 → carrier 被占用
Thread v1 = Thread.ofVirtual().start(() -> {
synchronized (MONITOR) {
try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}
});

// 正例:换成 ReentrantLock,阻塞时可安全卸载
Thread v2 = Thread.ofVirtual().start(() -> {
LOCK.lock();
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
LOCK.unlock();
}
});

v1.join(); v2.join();
System.out.println("诊断方式:启动参数加 -Djdk.tracePinnedThreads=short 可看到 synchronized 的钉扎堆栈");
}
}

禁忌二:ThreadLocal 滥用。 每个虚拟线程都有自己的 ThreadLocal 副本,百万虚拟线程 × 几个 ThreadLocal Map 会造成可观的内存膨胀。JDK 21 引入的 ScopedValue 正为虚拟线程设计的”结构化、不可变、有作用域”替代方案。应避免用 ThreadLocal 传递大对象,链路追踪等框架层尤其要评估开销。

禁忌三:CPU 密集任务几乎无收益。 虚拟线程不增加算力,纯计算时数量再多也只能占有限的核,反而增加调度与内存开销。CPU 密集仍应走 ForkJoinPool / parallelStream / 固定大小线程池。

禁忌四:不要池化虚拟线程。 虚拟线程创建极其廉价,”复用”反而带来状态泄漏与调试困难。直接用 newVirtualThreadPerTaskExecutor(),每任务一条新虚拟线程即可。

与响应式编程(Reactor / RxJava)的取舍:虚拟线程让”同步阻塞的写法”重新可行,可读性与可调试性远胜层层回调;但响应式在背压控制、流式窗口聚合、事件编排上仍有不可替代的优势。务实选择:新项目优先虚拟线程 + 结构化并发;已有成熟响应式栈不必迁移;网关、流式计算继续用响应式。

8.3 最佳实践清单

  1. 禁止使用 Executors 创建线程池,一律 new ThreadPoolExecutor(...) 显式声明七大参数。
  2. 队列必须有界,容量按”峰值 QPS × 可容忍等待时间”估算,宁可拒绝不可 OOM。
  3. 线程必须命名并配置 UncaughtExceptionHandler,否则线上 jstack 无法定位。
  4. 按业务隔离线程池(舱壁模式),慢服务不能共用核心池。
  5. 拒绝策略必须自定义:记录指标 + 落库待补偿 + 告警,绝不静默丢弃。
  6. 任务内部必须 try-catch,并重写 afterExecute 兜住异常(execute 提交的异常会被吞掉)。
  7. 参数以压测为准,公式只作起点;把 core/max/queue 接入配置中心支持动态调整。
  8. 监控四要素:活跃线程数、队列水位、拒绝次数、任务 P99 耗时,配水位告警。
  9. 优雅停机三段式:shutdown() → awaitTermination → 超时再 shutdownNow(),并钩子等待存量任务。
  10. IO 密集且高并发(JDK 21+)优先虚拟线程,synchronized 改 ReentrantLock,慎用 ThreadLocal,不要池化。

8.4 高频面试题速答

  1. 线程池的执行流程? 核心线程 → 队列 → 最大线程 → 拒绝策略;入队后还要 recheck 状态与线程数。
  2. 为什么先入队而不是先扩容线程? 创建线程昂贵,队列能缓冲说明核心线程足以消化,避免无谓扩容的切换与内存开销(Doug Lea 源码注释明示)。
  3. corePoolSize 与 maximumPoolSize 区别? 前者是常驻规模,后者是队列满后的上限;无界队列时后者失效。
  4. keepAliveTime 对核心线程生效吗? 默认不生效,需 allowCoreThreadTimeOut(true)。
  5. ctl 为什么用一个 int? 高 3 位状态 + 低 29 位线程数,使状态与数量的更新一次 CAS 完成;29 位上限约 5.37 亿。
  6. Worker 为什么继承 AQS? 实现不可重入独占锁,tryLock 成功即代表空闲,供 shutdown 精确中断空闲线程。
  7. shutdown 与 shutdownNow 区别? 前者平和:不接新任务、跑完队列、不中断运行中任务;后者激进:清空队列并返回、中断所有线程。
  8. 为什么禁止 Executors? 无界队列 OOM、Integer.MAX_VALUE 最大线程数导致线程爆炸、参数不可控、默认线程名无意义。
  9. submit 与 execute 区别? submit 返回 Future、异常被封装进 Future(get 时才抛);execute 无返回值、异常直接抛出并导致该线程重建。
  10. CPU 密集与 IO 密集线程数怎么设? N+1 与 2N(或 N × (1 + W/C)),最终以压测为准。
  11. CallerRunsPolicy 有什么用? 让提交者自己执行,形成反压降低提交速率;风险是占用上游线程、池关闭时静默丢弃。
  12. thenApply 与 thenCompose 区别? 前者同步转换为 CF<U>;后者接收返回 CF 的函数并扁平化为 CF<U>,避免嵌套。
  13. exceptionally/handle/whenComplete 差异? 分别对应:异常兜底、结果+异常统一转换、不改变结果的回调(类似 finally)。
  14. 为什么 supplyAsync 要指定线程池? 默认 ForkJoinPool.commonPool 全局共享、并行度 CPU-1、守护线程,IO 任务会饿死其它调用方。
  15. 工作窃取为何自己 LIFO、窃取 FIFO? LIFO 提升缓存局部性,FIFO 从队尾拿到更大的任务块,减少窃取次数与竞争。
  16. 什么时候不该用并行流? 数据量小、IO 密集、有状态操作、数据源难拆分。
  17. 虚拟线程适合什么场景? 高并发 IO 密集;不适合 CPU 密集;注意 synchronized 钉扎、ThreadLocal 膨胀、不要池化。
  18. 虚拟线程与响应式如何取舍? 虚拟线程胜在可读可调试,响应式胜在背压与流式编排。

到这里”如何用好线程”这条主线就完整了:从线程的成本认识池化的价值,从七大参数与决策链理解行为边界,从 ctl 位运算与 AQS Worker 看懂实现的精妙,从 CompletableFuture 把异步升级为编排,从 ForkJoin 认识另一种并行范式,最后用 虚拟线程看清这套体系的走向。下一篇将进入 JVM 层面,看看这些线程共享的内存是如何被划分、分配与回收的。