Java 基础体系 · 第 13/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java 线程与 Executor:生命周期、线程池、队列、Future 和取消
Java 并发代码通常同时包含两类对象:
- 任务:要执行的逻辑,例如实现了
Runnable或Callable<T>的对象; - 执行者:负责安排任务在哪个线程上执行,例如
Thread、Executor或ExecutorService。
Thread 解决的是“一个线程如何创建和运行”;Executor 解决的是“任务如何交给执行机制,以及如何管理执行资源”。线程池、工作队列、Future 和取消,都是围绕这两个层次展开的。
本文以 Java 25 为范围,区分 Java 规范保证、标准 API 契约和常见实现行为。
1. 从任务到线程:先区分三个生命周期
一个任务至少经历三个不同生命周期:
- 任务对象生命周期:对象被创建、提交、运行、结束;
- 线程生命周期:线程从
NEW进入可运行状态,最后变为TERMINATED; - Executor 生命周期:执行器接受任务、停止接收任务、等待工作线程结束。
这三个生命周期并不相同。
例如:
ExecutorService executor = Executors.newFixedThreadPool(2);
Future<Integer> future = executor.submit(() -> 42);
executor.shutdown(); // 不再接受新任务
int result = future.get(); // 已提交的任务仍可能正常完成
shutdown() 结束的是 Executor 的“接收任务”阶段,不是立即结束所有已经提交的任务。任务、工作线程和 Executor 的状态必须分别分析。
2. Thread 的生命周期和状态
2.1 创建线程并不等于启动线程
Thread thread = new Thread(() -> {
System.out.println("running");
});
System.out.println(thread.getState()); // NEW
thread.start(); // 请求 JVM 调度该线程
调用构造方法只创建了 Thread 对象;调用 start() 才使线程进入执行生命周期。
start() 有两个重要约束:
- 一个
Thread对象只能成功调用一次start(); - 第二次调用会抛出
IllegalThreadStateException。
Thread thread = new Thread(() -> {});
thread.start();
thread.start(); // IllegalThreadStateException
如果确实要再次执行相同逻辑,应该创建新的 Thread 对象,而不是复用已经结束的线程对象。
2.2 Java 线程状态
Thread.State 定义了六种观察状态:
| 状态 | 含义 |
|---|---|
NEW |
已创建但尚未调用 start() |
RUNNABLE |
正在运行,或已经可以运行,等待操作系统调度 |
BLOCKED |
等待进入 synchronized 监视器锁 |
WAITING |
无限期等待其他线程的特定动作 |
TIMED_WAITING |
带超时地等待 |
TERMINATED |
run() 执行结束,或因未捕获异常退出 |
这些状态是 Java 对线程状态的抽象,并不等于操作系统线程的所有底层状态。
例如,RUNNABLE 同时覆盖:
- 线程正在 CPU 上执行;
- 线程已经就绪但暂时没有获得 CPU;
- 某些 JVM 实现中正在执行本地代码或等待底层资源。
因此,看到大量线程处于 RUNNABLE,不能直接推断它们都在消耗 CPU;应结合 CPU 使用率、线程栈和锁信息判断。
2.3 状态转换
典型状态路径如下:
stateDiagram-v2
[*] --> NEW
NEW --> RUNNABLE: start()
RUNNABLE --> BLOCKED: 等待 synchronized 锁
BLOCKED --> RUNNABLE: 获得锁
RUNNABLE --> WAITING: wait()/join()/park()
RUNNABLE --> TIMED_WAITING: sleep()/wait(timeout)/join(timeout)
WAITING --> RUNNABLE: 被唤醒或许可可用
TIMED_WAITING --> RUNNABLE: 超时、唤醒或许可可用
RUNNABLE --> TERMINATED: run() 返回
RUNNABLE --> TERMINATED: 未捕获异常
状态转换中的几个常见误解:
Thread.sleep()不会释放已经持有的锁;Object.wait()会释放对应对象的监视器锁,并进入等待;join()等待的是另一个线程结束;interrupt()不是强制杀死线程,而是改变中断状态,或唤醒部分可中断阻塞操作。
3. 线程运行、异常和中断
3.1 run() 与 start() 完全不同
Thread thread = new Thread(() -> {
System.out.println(Thread.currentThread().getName());
});
thread.run(); // 普通方法调用,在当前线程执行
thread.start(); // 创建并调度新的线程执行
如果当前线程名是 main,直接调用 run() 时输出仍然是 main。只有 start() 才建立新的线程执行路径。
3.2 未捕获异常不会自动传给提交者
Thread thread = new Thread(() -> {
throw new RuntimeException("boom");
});
thread.setUncaughtExceptionHandler((t, e) ->
System.err.println(t.getName() + ": " + e));
thread.start();
线程中的未捕获异常会终止该线程,并交给未捕获异常处理器;它不会像普通方法返回值一样自动传播到创建线程的调用栈。
在线程池中需要额外区分 execute 和 submit,后文会说明:submit 通常把异常保存到 Future 中,而不是立即交给线程的未捕获异常处理器。
3.3 中断是协作式取消信号
Thread.interrupt() 的语义取决于目标线程当前做什么:
- 如果目标线程正在
sleep、wait、join等可中断阻塞中,通常会收到InterruptedException,并清除中断标志; - 如果目标线程没有处于这些阻塞调用中,中断标志会被设置;
- 如果目标线程完全忽略中断,线程不会因此自动停止。
Thread thread = new Thread(() -> {
try {
while (!Thread.currentThread().isInterrupted()) {
doSmallUnitOfWork();
Thread.sleep(100);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // 恢复中断状态
} finally {
cleanup();
}
});
thread.start();
thread.interrupt();
这里的控制流是:
- 循环通过
isInterrupted()检查中断; sleep()可能抛出InterruptedException;- 捕获异常后恢复中断状态;
finally执行清理;- 线程最终返回。
如果捕获 InterruptedException 后既不重新抛出,也不恢复中断状态,调用方可能无法继续观察到取消信号:
try {
Thread.sleep(1000);
} catch (InterruptedException ignored) {
// 错误示例:吞掉中断
}
正确处理方式通常是:
- 方法能够声明异常时,继续抛出;
- 不能声明时,调用
Thread.currentThread().interrupt(),然后结束当前任务或进入明确的恢复逻辑。
3.4 中断与内存可见性
中断状态不是适合所有取消协议的唯一共享变量。若任务使用独立的取消标志,该标志必须具有可见性:
final class CancellableTask implements Runnable {
private volatile boolean cancelled;
void cancel() {
cancelled = true;
}
@Override
public void run() {
while (!cancelled) {
doSmallUnitOfWork();
}
}
private void doSmallUnitOfWork() {
// 必须是有限时间内能够返回的工作
}
}
volatile 保证对该变量的写入对后续读取可见,但不保证复合操作的原子性。例如:
if (!cancelled) {
cancelled = true;
}
这不是一个有意义的原子“检查并设置”操作。需要原子条件更新时,应考虑 AtomicBoolean.compareAndSet 或锁。
4. Java 内存模型:为什么提交和 Future.get() 能传递结果
线程池正确工作的基础不仅是队列,还包括线程之间的内存可见性。
Java 内存模型中的 happens-before 是一种可传递的先行关系。若动作 A happens-before 动作 B,则 A 的效果对 B 可见,并且执行顺序不能被重排成违反该关系的结果。
常见规则包括:
- 对一个
volatile变量的写,happens-before 后续对该变量的读; - 解锁一个监视器 happens-before 随后对同一个监视器的加锁;
- 线程中的动作 happens-before 另一个线程成功从
Thread.join()返回; - 根据
ExecutorService的契约,提交任务之前的动作 happens-before 任务开始执行; - 任务中的动作 happens-before 另一个线程成功从该任务的
Future.get()返回。
因此下面的代码有明确的可见性保证:
ExecutorService executor = Executors.newSingleThreadExecutor();
int[] input = {10};
Future<Integer> future = executor.submit(() -> input[0] * 2);
int result = future.get(); // result 为 20
executor.shutdown();
input[0] = 10 发生在提交前,任务读取 input[0] 时可以观察到该写入;任务计算出的结果通过 Future.get() 返回给调用线程。
但这不意味着所有共享对象都自动变成线程安全对象。若任务在执行期间和其他线程并发修改同一个可变对象,仍然需要锁、volatile、原子类或其他同步协议。
4.1 一个错误的共享状态例子
class Counter {
int value;
void increment() {
value++;
}
}
value++ 实际上包含读取、加一、写回三个动作。两个线程可能读取同一个旧值,导致更新丢失。
如果只需要原子递增,可以使用:
AtomicInteger counter = new AtomicInteger();
executor.submit(counter::incrementAndGet);
如果要维护多个字段之间的一致性,单个 AtomicInteger 通常不够,可能需要使用锁保护整个不变量。
5. Executor:把任务提交与线程管理分离
5.1 Executor 的最小抽象
@FunctionalInterface
public interface Executor {
void execute(Runnable command);
}
调用者只需要提交 Runnable:
Executor executor = command -> {
Thread thread = new Thread(command);
thread.start();
};
executor.execute(() -> System.out.println("running"));
调用者不再直接管理线程对象,但这个简单实现每提交一个任务就创建一个平台线程,无法限制并发度,也没有统一关闭和结果管理能力。
5.2 ExecutorService 增加生命周期和结果管理
ExecutorService 在 Executor 之上增加了:
submit:提交Runnable或Callable<T>;Future:观察任务结果、异常和取消状态;shutdown:停止接收新任务;shutdownNow:尝试中断执行中的任务,并返回尚未开始的任务;awaitTermination:等待 Executor 终止;invokeAll、invokeAny:批量执行任务。
典型生命周期:
stateDiagram-v2
[*] --> RUNNING
RUNNING --> SHUTDOWN: shutdown()
RUNNING --> STOP: shutdownNow()
SHUTDOWN --> TIDYING: 已提交任务全部完成
STOP --> TIDYING: 工作线程结束
TIDYING --> TERMINATED: terminated()
ExecutorService 的 isShutdown() 和 isTerminated() 含义不同:
isShutdown():已经不再接受新任务;isTerminated():关闭流程已完成,工作线程已退出。
调用 shutdown() 后,isShutdown() 可能立刻为 true,但 isTerminated() 仍然为 false。
6. 线程池的核心组成
典型平台线程池包含:
- 工作线程:真正执行任务;
- 工作队列:保存暂时没有线程执行的任务;
- 线程工厂:创建和命名线程;
- 拒绝策略:线程和队列都无法接收任务时的处理方式;
- 生命周期控制:关闭、等待和中断。
ThreadPoolExecutor 的关键参数是:
ThreadPoolExecutor(
int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler
)
6.1 提交流程
对于默认配置下的 ThreadPoolExecutor.execute(command),核心判断顺序可以抽象为:
1. 当前工作线程数 < corePoolSize?
是:创建新的核心线程执行任务。
否:进入第 2 步。
2. workQueue.offer(command) 成功?
是:任务进入队列。
否:进入第 3 步。
3. 当前工作线程数 < maximumPoolSize?
是:创建非核心线程执行任务。
否:执行拒绝策略。
这个顺序非常重要。很多人以为“任务多了就扩容到 maximumPoolSize”,但如果工作队列是无界队列,第二步几乎永远成功,线程池通常不会扩展到 maximumPoolSize。
例如:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
2,
4,
30,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>()
);
在常见实现和该队列行为下:
- 前两个任务创建核心线程;
- 后续任务进入无界队列;
maximumPoolSize = 4基本不会触发;- 如果生产速度长期高于消费速度,队列可能持续增长,最终造成内存压力。
因此 maximumPoolSize 只有在队列拒绝继续接收任务时才有机会发挥作用。
7. 工作队列决定线程池的行为
7.1 SynchronousQueue
SynchronousQueue 不保存元素,每次 put 都必须直接交给一个正在等待的消费者。
new ThreadPoolExecutor(
2,
8,
30,
TimeUnit.SECONDS,
new SynchronousQueue<>()
);
这种配置倾向于:
- 任务没有排队缓冲;
- 没有空闲工作线程时尝试创建新线程;
- 达到最大线程数后快速触发拒绝。
它适合需要快速暴露过载、并且允许并发线程数受 maximumPoolSize 限制的场景,但不适合无限制地创建线程。
7.2 有界队列
new ArrayBlockingQueue<>(100)
有界队列同时限制了:
- 工作线程数量;
- 队列中等待任务的数量。
当核心线程用尽、队列容量用尽、线程数又达到最大值时,任务会被拒绝。这个拒绝是重要的背压信号,而不是一个应该被无条件吞掉的异常。
7.3 无界队列
new LinkedBlockingQueue<>()
无界队列通常意味着排队容量没有业务上限。即使底层实现对内部容量使用了一个很大的最大值,也不能把它当成稳定的业务容量控制。
风险是:
到达速率 > 处理速率
=> 队列长度增长
=> 任务等待时间增长
=> 堆内存增长
=> GC 压力或 OutOfMemoryError
无界队列可能减少拒绝,但把过载转化为延迟和内存风险。
7.4 队列、公平性和顺序
常见阻塞队列通常提供先进先出语义,但线程池整体不等于严格的全局 FIFO 系统:
- 多个工作线程同时执行任务,完成顺序可能不同;
- 任务入队和线程取任务之间存在并发竞争;
- 不同优先级、不同任务耗时会造成完成顺序差异;
- 使用
PriorityBlockingQueue时,顺序由优先级决定,但必须处理优先级反转和长期饥饿问题。
队列中的 FIFO 只能说明“从队列取出时的相对顺序”,不能保证业务完成顺序。
8. 创建线程池时不要忽略线程工厂和拒绝策略
8.1 命名线程并配置异常处理
AtomicInteger sequence = new AtomicInteger();
ThreadFactory factory = task -> {
Thread thread = new Thread(task);
thread.setName("image-worker-" + sequence.incrementAndGet());
thread.setUncaughtExceptionHandler((t, e) ->
System.err.println(t.getName() + " failed: " + e));
return thread;
};
线程命名让线程转储、监控和日志能够关联到具体执行器。线程工厂还可以设置:
- 是否为守护线程;
- 线程优先级;
- 未捕获异常处理器;
- 线程组等属性。
守护线程不会阻止 JVM 退出,因此不能把必须完成的持久化、事务提交或资源释放任务仅交给守护线程。
8.2 拒绝策略
当线程池无法接受新任务时,内置策略包括:
| 策略 | 行为 |
|---|---|
AbortPolicy |
抛出 RejectedExecutionException |
CallerRunsPolicy |
由提交任务的调用线程执行任务 |
DiscardPolicy |
静默丢弃任务 |
DiscardOldestPolicy |
丢弃队列头部任务后重试提交 |
默认是 AbortPolicy。
CallerRunsPolicy 能形成一种简单背压:提交者自己执行任务,因此提交速度被处理速度拖慢。但如果提交者是 HTTP 请求线程、事件循环线程或关键调度线程,可能把耗时任务带入不应阻塞的线程。
DiscardPolicy 和 DiscardOldestPolicy 只有在业务明确允许丢失任务时才可使用。否则,静默丢弃会让故障变成数据缺失。
9. execute、submit 和 Future
9.1 execute:只提交 Runnable
executor.execute(() -> {
System.out.println("task");
});
它没有结果对象。任务抛出的未捕获异常通常会导致执行该任务的工作线程结束,并交给线程的异常处理机制。
9.2 submit:返回 Future
Future<Integer> future = executor.submit(() -> 21 * 2);
try {
Integer result = future.get();
System.out.println(result); // 42
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} catch (ExecutionException e) {
Throwable cause = e.getCause();
cause.printStackTrace();
}
submit 将任务包装为可跟踪的计算。Future.get() 可能出现:
InterruptedException:等待结果的当前线程被中断;ExecutionException:任务执行失败,真正异常在getCause()中;CancellationException:任务已取消;TimeoutException:带超时的get在期限内未完成。
任务异常不会直接从 submit 调用点抛出:
Future<?> future = executor.submit(() -> {
throw new IllegalStateException("failed");
});
System.out.println("submit returned"); // 通常会输出
future.get(); // 此处抛 ExecutionException
如果只调用 submit 而永远不检查 Future,任务失败可能被悄悄隐藏。
9.3 Future 表示什么
Future<T> 表示一个可能尚未完成的计算,而不是结果本身。它支持:
boolean isDone();
boolean isCancelled();
T get() throws InterruptedException, ExecutionException;
T get(long timeout, TimeUnit unit)
throws InterruptedException, ExecutionException, TimeoutException;
boolean cancel(boolean mayInterruptIfRunning);
状态可以抽象为:
未完成
├─ 正常完成 -> get 返回结果
├─ 执行失败 -> get 抛 ExecutionException
└─ 被取消 -> get 抛 CancellationException
isDone() 在正常完成、失败和取消三种情况下都可能返回 true,它不代表任务成功。
10. Future.get、超时和调用线程中断
10.1 无限等待可能造成线程堆积
Result result = future.get();
如果任务依赖一个永远不返回的网络调用,等待 get() 的线程也会一直占用资源。更常见的做法是设置超时:
try {
Result result = future.get(500, TimeUnit.MILLISECONDS);
} catch (TimeoutException e) {
future.cancel(true);
throw new RuntimeException("task timed out", e);
} catch (InterruptedException e) {
future.cancel(true);
Thread.currentThread().interrupt();
throw new RuntimeException("caller interrupted", e);
} catch (ExecutionException e) {
throw new RuntimeException("task failed", e.getCause());
}
超时只表示“调用者不再愿意等待”,不等于任务已经停止。必须显式取消,且任务本身必须响应取消。
10.2 超时与取消的竞态
以下事件可能并发发生:
调用者 get 超时
任务刚好正常完成
调用者调用 cancel(true)
因此不能假设 cancel(true) 一定返回 true。如果任务已经完成,取消会失败;如果任务尚未完成,取消可能成功。业务代码必须接受这种竞态,并以 Future 的最终状态为准。
11. Future 的取消:停止等待,不等于停止执行
调用:
future.cancel(true);
true 的含义是:如果任务已经开始执行,尝试中断执行线程。
它不保证:
- 线程立即终止;
- 正在执行的本地方法被强制打断;
- 阻塞在不可中断 I/O 上的任务立即返回;
- 外部副作用自动回滚;
- 已经发送的消息或已提交的数据库事务被撤销。
如果任务尚未开始执行,取消通常可以让任务不再执行;如果任务已经完成,取消失败;如果任务正在运行,取消成功只代表 Future 被标记为取消,并发出中断请求。
11.1 可取消任务的完整例子
import java.util.concurrent.*;
public class CancellationDemo {
public static void main(String[] args)
throws InterruptedException {
ExecutorService executor = Executors.newSingleThreadExecutor();
Future<?> future = executor.submit(() -> {
try {
while (true) {
System.out.println("processing");
Thread.sleep(200);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
System.out.println("task cancelled");
} finally {
System.out.println("cleanup");
}
});
Thread.sleep(700);
boolean cancelled = future.cancel(true);
System.out.println("cancelled = " + cancelled);
executor.shutdown();
if (!executor.awaitTermination(2, TimeUnit.SECONDS)) {
executor.shutdownNow();
}
}
}
执行过程通常是:
- 单线程池创建一个工作线程;
- 任务循环运行;
- 主线程调用
cancel(true); - 工作线程在
sleep中收到InterruptedException; - 任务恢复中断状态并退出;
finally执行清理;- Executor 关闭。
11.2 不响应取消的任务
executor.submit(() -> {
while (true) {
calculate();
}
});
如果 calculate() 不检查中断,也不调用可中断阻塞操作,那么:
future.cancel(true);
通常只能使 Future 进入取消状态,而任务仍可能继续消耗 CPU。
这也是 shutdownNow() 不能保证立即停止所有任务的根本原因:它的主要机制也是中断工作线程,而中断是协作式的。
12. shutdown、shutdownNow 和资源释放
12.1 正常关闭
executor.shutdown();
try {
if (!executor.awaitTermination(30, TimeUnit.SECONDS)) {
executor.shutdownNow();
}
} catch (InterruptedException e) {
executor.shutdownNow();
Thread.currentThread().interrupt();
}
这个关闭流程的含义是:
shutdown():不再接受新任务;- 已提交任务继续执行;
awaitTermination等待工作线程结束;- 超时后调用
shutdownNow()尝试中断剩余任务; - 如果等待线程自身被中断,也应继续关闭并恢复中断状态。
12.2 强制关闭
List<Runnable> notStarted = executor.shutdownNow();
shutdownNow() 通常会:
- 停止接受新任务;
- 尝试中断正在执行的任务;
- 返回队列中尚未开始执行的任务。
返回列表不包括已经开始执行的任务,也不能说明被中断的任务是否真的停止。
任务可能在以下时刻与关闭并发:
shutdownNow() 检查队列
|
|-- 工作线程刚好取走某任务
|
返回“未开始任务”列表
所以返回结果只能作为未启动任务的近似管理依据,不能替代任务自身的幂等和恢复设计。
12.3 try-with-resources
ExecutorService 实现了 AutoCloseable,可以写成:
try (ExecutorService executor = Executors.newFixedThreadPool(2)) {
Future<Integer> future = executor.submit(() -> 42);
System.out.println(future.get());
}
退出 try 块时会关闭 Executor。此方式适合 Executor 生命周期明确属于当前代码块的场景。若 Executor 是应用级共享组件,不应让某个局部方法随意关闭它。
13. Executors 工厂与显式配置
常见工厂方法包括:
Executors.newFixedThreadPool(4);
Executors.newSingleThreadExecutor();
Executors.newCachedThreadPool();
Executors.newScheduledThreadPool(2);
Executors.newVirtualThreadPerTaskExecutor();
13.1 固定线程池
newFixedThreadPool(n) 通常使用固定数量的工作线程和无界队列。它限制了并发执行线程数,但不限制等待任务数量。
适用于任务量可控、允许排队的场景;如果任务持续积压,需要额外监控队列长度和等待时间。
13.2 单线程 Executor
newSingleThreadExecutor() 保证同一时刻只有一个任务执行。它适合串行化某类操作,但不应被误认为是通用同步机制:
- 其他线程仍可能直接访问共享状态;
- 任务内部抛出异常不会让业务自动恢复;
- 无界排队仍可能导致任务积压。
13.3 缓存线程池
newCachedThreadPool() 倾向于为突发任务创建线程,并回收一段时间未使用的线程。它没有一个由调用者传入的固定最大线程数,因此面对无限制提交或慢任务时可能创建大量平台线程。
13.4 显式构造的价值
当队列容量、线程命名、最大线程数和拒绝行为属于业务约束时,直接构造 ThreadPoolExecutor 更容易审查:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
4,
8,
30,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(200),
factory,
new ThreadPoolExecutor.CallerRunsPolicy()
);
这组参数的实际含义是:
- 最多同时有 8 个工作线程;
- 常态保留 4 个核心线程;
- 最多缓存 200 个等待任务;
- 队列和线程都满时由提交者执行任务;
- 非核心线程空闲 30 秒后可以回收。
不能仅凭“核心线程数等于 CPU 核心数”推导出最优配置。CPU 密集型任务、阻塞 I/O 任务、混合任务的等待和计算比例不同,所需并发度也不同。
14. ScheduledExecutorService:延迟与周期任务
ScheduledExecutorService 用于延迟执行或周期执行:
ScheduledExecutorService scheduler =
Executors.newScheduledThreadPool(1);
ScheduledFuture<?> handle = scheduler.scheduleAtFixedRate(
() -> System.out.println(System.currentTimeMillis()),
0,
1,
TimeUnit.SECONDS
);
两种周期方式不同:
scheduleAtFixedRate:以固定速率安排,目标时间基于初始时间加周期;scheduleWithFixedDelay:一次执行结束后,再等待指定延迟执行下一次。
如果周期任务执行时间长于周期,固定速率不会并发执行同一个周期任务;后续执行可能延迟。周期任务抛出未捕获异常时,后续周期执行可能停止,因此应在任务内部处理预期异常。
取消周期任务:
handle.cancel(false);
scheduler.shutdown();
周期调度器适合轻量调度,不适合替代可靠的分布式定时系统。进程崩溃、机器重启和时钟变化都可能影响它。
15. 批量任务:invokeAll 和 invokeAny
15.1 invokeAll
List<Callable<Integer>> tasks = List.of(
() -> 1,
() -> 2,
() -> 3
);
List<Future<Integer>> futures = executor.invokeAll(tasks);
for (Future<Integer> future : futures) {
try {
System.out.println(future.get());
} catch (ExecutionException e) {
System.err.println("one task failed: " + e.getCause());
}
}
invokeAll 等待所有任务完成,并按输入任务顺序返回对应的 Future 列表,而不是按完成顺序返回。
带超时时,尚未完成的任务会被取消:
List<Future<Integer>> futures =
executor.invokeAll(tasks, 1, TimeUnit.SECONDS);
调用线程被中断时,批量等待会抛出 InterruptedException,调用方仍应恢复或传递中断。
15.2 invokeAny
Integer result = executor.invokeAny(tasks);
它返回一个成功完成任务的结果,并取消其他尚未完成的任务。这里的“成功”意味着没有抛异常;如果所有任务都失败,调用会报告失败。
invokeAny 适合多个等价副本竞速,但取消其他任务仍然是协作式的,不能假设后续任务立刻停止。
16. 平台线程池与虚拟线程 Executor
Java 25 中,虚拟线程已经是标准能力。可以使用:
try (ExecutorService executor =
Executors.newVirtualThreadPerTaskExecutor()) {
Future<String> future = executor.submit(() -> {
Thread.sleep(100);
return "ok";
});
System.out.println(future.get());
}
这个 Executor 的语义与固定平台线程池不同:
- 每个提交的任务通常对应一个新的虚拟线程;
- 它不是通过复用少量工作线程来限制任务并发;
close()会关闭 Executor,并等待其任务结束;- 任务数量仍应受到业务容量、下游连接数、内存和速率限制约束。
“虚拟线程便宜”不等于“可以无限提交任务”。如果同时向数据库提交一百万个查询,瓶颈仍然是数据库连接、数据库 CPU 和应用内存。
16.1 虚拟线程不是传统线程池大小替代品
平台线程池通常表达:
最多同时执行 N 个任务
虚拟线程 Executor 更接近:
为每个任务提供一个轻量线程
如果需要限制并发访问下游服务,应使用显式容量控制,例如 Semaphore:
Semaphore permits = new Semaphore(50);
try (ExecutorService executor =
Executors.newVirtualThreadPerTaskExecutor()) {
Future<?> future = executor.submit(() -> {
permits.acquireUninterruptibly();
try {
callDownstream();
} finally {
permits.release();
}
});
future.get();
}
这里限制的是 callDownstream() 的并发数量,而不是虚拟线程总数。
16.2 Pinning 的边界
虚拟线程运行 Java 代码时,会被调度到少量平台线程上执行。某些操作可能使虚拟线程暂时固定在承载它的平台线程上,这称为 pinning。长时间 pinning 会降低调度并行能力,尤其是:
- 在
synchronized临界区内执行长时间阻塞操作; - 执行某些本地方法或受底层实现影响的操作。
因此,虚拟线程适合大量阻塞式 I/O,但仍应避免在锁保护区域内进行慢 I/O:
synchronized (lock) {
remoteCall(); // 风险:临界区过长,且可能发生 pinning
}
可以缩小锁范围,或使用适合场景的 ReentrantLock,但不能机械地把所有 synchronized 替换掉;锁仍然是保护不变量的重要工具。
虚拟线程的调度、pinning 诊断和结构化并发属于 Java 25 并发体系中的相关主题,但它们不改变本文关于任务、队列、Future 和协作式取消的基本契约。
17. 常见失败模式和诊断路径
17.1 线程池“卡住”
可能的因果链:
核心线程正在等待外部服务
=> 队列继续增长
=> 新任务无法及时执行
=> 调用方 get 超时
=> 更多重试任务进入队列
=> 过载进一步加剧
诊断应同时观察:
- 活跃线程数;
- 当前池大小;
- 队列长度;
- 已完成任务数;
- 任务平均和最大执行时间;
Future.get超时数量;- 下游服务响应时间;
- 线程转储中的等待位置。
可以通过 ThreadPoolExecutor 的监控方法获取基础指标:
ThreadPoolExecutor executor = ...;
System.out.println("pool size = " + executor.getPoolSize());
System.out.println("active = " + executor.getActiveCount());
System.out.println("queue = " + executor.getQueue().size());
System.out.println("completed = " + executor.getCompletedTaskCount());
这些指标是瞬时观察值,不能单独证明不存在竞态或隐藏积压,但能帮助确认“线程不足”还是“下游阻塞”。
17.2 Future.get() 全部超时
不能只增加超时时间。应区分:
- 任务本身耗时过长;
- 任务卡在锁竞争;
- 任务卡在 I/O;
- Executor 队列等待时间过长;
- 任务已经被取消;
- 调用线程本身被中断。
可以把“排队时间”和“执行时间”分别记录:
long submittedAt = System.nanoTime();
Future<?> future = executor.submit(() -> {
long startedAt = System.nanoTime();
recordQueueDelay(startedAt - submittedAt);
doWork();
recordExecutionTime(System.nanoTime() - startedAt);
});
这段代码能区分任务是“没拿到线程”还是“拿到线程后执行慢”。生产环境中应避免把共享的可变计时状态无保护地写入。
17.3 任务异常消失
错误代码:
executor.submit(() -> {
parseInput(); // 可能抛异常
});
如果不保存和检查返回的 Future,异常可能只存在于内部状态中。
改进方式:
Future<?> future = executor.submit(() -> parseInput());
try {
future.get();
} catch (ExecutionException e) {
log.error("task failed", e.getCause());
}
对于不需要结果但必须检测失败的任务,也应保留 Future,或者封装统一的提交方法,在任务内部记录异常。
17.4 关闭顺序错误
错误顺序:
executor.shutdown();
executor.submit(task); // RejectedExecutionException
关闭是单向生命周期变化。应用停止时通常应:
- 停止产生新任务;
- 调用
shutdown(); - 等待有限时间;
- 必要时调用
shutdownNow(); - 记录仍未结束的任务;
- 恢复被中断的关闭线程状态。
如果任务会向另一个 Executor 提交子任务,必须考虑关闭顺序和相互等待,否则可能发生关闭阶段死锁或永久等待。
18. 任务设计:有限步骤、幂等和资源边界
线程池只能管理线程和任务调度,不能替业务决定副作用如何恢复。
一个可取消任务通常应满足:
重复检查取消信号
-> 将工作拆成有限时间的小步骤
-> 每个外部副作用具有明确提交点
-> 失败和取消路径执行清理
-> 重试时保证幂等或使用去重标识
例如,文件上传任务可能已经完成远端写入,但本地任务被取消。Future 进入取消状态并不代表远端文件自动删除。要实现业务取消,需要额外的协议:
- 上传接口支持取消;
- 记录上传会话并执行补偿删除;
- 使用幂等请求标识;
- 明确“取消成功”只代表不再继续,而不是回滚已经完成的副作用。
线程中断和业务回滚是两个不同层次的问题。
19. 如何选择执行模型
可以按任务特征推导,而不是只按线程数量选择:
CPU 密集型
任务主要消耗 CPU,增加线程通常只会增加上下文切换。应控制并发度,避免让大量任务同时争用处理器。
阻塞 I/O 型
任务主要等待网络、文件或数据库。平台线程池可以通过增加线程隐藏等待时间,但线程数量过大有栈内存、调度和下游过载风险;虚拟线程可以降低每个阻塞任务的线程资源成本,但仍需限制下游并发。
混合型
如果任务先计算、再 I/O、再计算,单一线程池可能把不同阶段互相拖慢。可以拆分阶段,但拆分后要重新处理队列容量、异常传播、取消传播和结果关联。
有明确容量上限的系统
使用有界队列或显式信号量,把超载转化为可观察的拒绝、降级或排队,而不是让无界队列隐藏问题。
20. 一个完整的可运行示例
下面的程序展示:
- 有界线程池;
- 命名线程;
Callable返回结果;- 任务异常通过
ExecutionException传播; - 超时后的协作式取消;
- 正确关闭 Executor。
import java.util.List;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class ExecutorEndToEnd {
public static void main(String[] args) {
AtomicInteger id = new AtomicInteger();
ThreadFactory factory = task -> {
Thread t = new Thread(task);
t.setName("worker-" + id.incrementAndGet());
return t;
};
ExecutorService executor = new ThreadPoolExecutor(
2,
4,
10,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(10),
factory,
new ThreadPoolExecutor.AbortPolicy()
);
try {
Future<String> success = executor.submit(() -> {
Thread.sleep(100);
return "success";
});
Future<String> failure = executor.submit(() -> {
throw new IllegalArgumentException("bad input");
});
Future<String> slow = executor.submit(() -> {
try {
Thread.sleep(10_000);
return "too late";
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return "cancelled cooperatively";
}
});
System.out.println(success.get());
try {
failure.get();
} catch (ExecutionException e) {
System.out.println("failure: " + e.getCause());
}
try {
System.out.println(slow.get(200, TimeUnit.MILLISECONDS));
} catch (TimeoutException e) {
boolean accepted = slow.cancel(true);
System.out.println("cancel requested: " + accepted);
} catch (InterruptedException e) {
slow.cancel(true);
Thread.currentThread().interrupt();
} catch (ExecutionException e) {
System.out.println("slow task failed: " + e.getCause());
}
} catch (InterruptedException e) {
executor.shutdownNow();
Thread.currentThread().interrupt();
} catch (ExecutionException e) {
executor.shutdownNow();
throw new RuntimeException(e.getCause());
} finally {
executor.shutdown();
try {
if (!executor.awaitTermination(2, TimeUnit.SECONDS)) {
List<Runnable> notStarted = executor.shutdownNow();
System.out.println("not started: " + notStarted.size());
}
} catch (InterruptedException e) {
executor.shutdownNow();
Thread.currentThread().interrupt();
}
}
}
}
预期现象包括:
success
failure: java.lang.IllegalArgumentException: bad input
cancel requested: true
慢任务是否已经输出 "cancelled cooperatively",取决于取消请求与工作线程执行到 sleep 的时序,但它正确处理了中断,因此通常可以较快结束。cancel(true) 返回 true 只表示取消请求被 Future 接受,不应解释为“任务已完成清理”。
21. 最容易混淆的几个结论
Thread.run()不创建新线程,Thread.start()才会;RUNNABLE不等于“正在占用 CPU”;shutdown()不会取消已提交任务;shutdownNow()只是尝试中断,并不保证任务停止;Future.cancel(true)不等于强制终止;Future.isDone()不等于执行成功;- 无界队列会削弱
maximumPoolSize的作用; submit中的异常通常要通过Future.get()观察;- 虚拟线程降低线程成本,但不会消除数据库、网络和内存容量限制;
- 中断是线程级取消信号,业务副作用的回滚必须由业务协议实现;
- 队列顺序不保证任务完成顺序;
- 线程安全不由 Executor 自动提供,共享状态仍需遵守 Java 内存模型和同步规则。
理解这些边界后,Thread、线程池、队列和 Future 不再是孤立 API,而可以被看作一条完整的数据流:
任务创建
-> 提交到 Executor
-> 工作线程取出或进入队列
-> 任务执行
-> 正常结果 / 异常 / 取消
-> Future 观察结果
-> Executor 有序关闭
任何一步缺少容量、可见性、异常传播或取消协议,最终都可能表现为线程泄漏、任务丢失、队列堆积、请求超时或关闭不完整。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java 网络与 HTTP Client:连接、TLS、超时、流和错误恢复
- 下一篇:Java 25 虚拟线程与结构化并发:调度、Pinning、取消和容量
- 延伸:Java 内存模型与同步:happens-before、锁、volatile、Atomic 和 VarHandle
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论