ExecutorService
本章以企业项目最常见的 JDK 7/JDK 8 线程池模型 为主线讲解:
Executor、ExecutorService、ThreadPoolExecutor、Callable、Future。文中涉及 JDK 19/21 之后新增的Future默认方法时,只作为版本演进补充,不作为基础学习前提。
ExecutorService 是线程池体系里的核心执行器接口,继承自 Executor,常见实现是 ThreadPoolExecutor。它提供了两类常用提交方式:不关心返回值时用 execute(),需要结果、异常或取消时用 submit() 返回 Future。
为什么需要 ExecutorService
如果每来一个任务都 new Thread(),短期看起来简单,但高并发时会出现三个问题:线程创建销毁成本高、线程数量不可控、任务失败和排队策略不可管理。ExecutorService 把“任务提交”和“线程执行”分开,让线程可以复用,也让队列、拒绝策略、关闭流程变成可控配置。
flowchart TD
A["业务提交任务"] --> B["ExecutorService"]
B --> C{"是否有可用线程"}
C -- "有" --> D["复用工作线程执行"]
C -- "没有" --> E["进入任务队列或触发拒绝策略"]
D --> F["任务完成后线程继续复用"]如果不会线程池原理,最容易把线程池当成“无限异步工具”。结果是队列堆满、内存升高、请求延迟变长,最后触发拒绝策略或服务整体不可用。
execute方法
execute() 方法是Executor接口定义的方法,当你只需要提交一个简单的Runnable任务不关心其返回值时可以使用这个方法。定义如下:
public interface Executor {
/**
* Executes the given command at some time in the future. The command
* may execute in a new thread, in a pooled thread, or in the calling
* thread, at the discretion of the {@code Executor} implementation.
* 在将来的某个时间执行给定的命令。该命令可以在新线程、池线程或调用线程中执行,具体由实现决定 Executor 。
*
* @param command the runnable task
* command – Runnable 任务
*
* @throws RejectedExecutionException if this task cannot be
* accepted for execution
* RejectedExecutionException – 如果无法接受此任务执行
*
* @throws NullPointerException if command is null
* NullPointerException – 如果 command 为 null
*/
void execute(Runnable command);
}submit方法
submit()是ExecutorService接口中定义的抽象方法用于提交一个任务,并返回一个表示该任务的 Future 对象。Future 对象可以用来检查任务是否完成、等待任务完成并获取其结果。抽象方法如下:
/** 提交一个返回值的任务以供执行,并返回一个 Future,表示该任务的待处理结果。Future 的方法 get 将在成功完成后返回任务的结果。
* 如果您想立即阻止等待任务,您可以使用表单 result = exec.submit(aCallable).get();
* 注意: 该 Executors 类包含一组方法,这些方法可以将其他一些常见的类似闭包的对象 java.security.PrivilegedAction Callable 转换为 form,以便可以提交它们。
* 参数:task – 要提交的任务
* 返回:一个 Future 表示任务的待完成
* 抛出:RejectedExecutionException – 如果无法安排任务执行,NullPointerException – 如果任务为 null
**/
<T> Future<T> submit(Callable<T> task);
/**提交一个 Runnable 任务以供执行,并返回一个表示该任务的 Future。Future 的方法 get 将在成功完成后返回给定的结果。
* 参数:task – 要提交的任务 result – 要返回的结果
* 返回:一个 Future 表示任务的待完成
* 抛出:RejectedExecutionException – 如果无法安排任务执行,NullPointerException – 如果任务为 null
**/
<T> Future<T> submit(Runnable task, T result);
/**提交一个 Runnable 任务以供执行,并返回一个表示该任务的 Future。Future 的方法get将在成功完成后返回null。
* 参数:task – 要提交的任务
* 返回:一个 Future 表示任务的待完成
* 抛出:RejectedExecutionException – 如果无法安排任务执行,NullPointerException – 如果任务为 null
**/
Future<?> submit(Runnable task);Future相关
上述submit方法执行返回一个Future对象,其常见抽象方法如下:
/**
* 尝试取消此任务的执行。如果任务已经完成或取消,或者由于其他原因无法取消,则此方法无效。否则,如果此任务在调用时cancel尚未启动,则此任务不应运行。
* 如果任务 已经启动,则参数 mayInterruptIfRunning 确定执行此任务的线程(当实现知道时)是否在尝试停止任务时被中断。
* 此方法的返回值不一定指示任务现在是否已取消;使用 isCancelled.
* 参数:mayInterruptIfRunning – true 如果执行此任务的线程应中断(如果实现已知该线程);否则,允许正在进行的任务完成
* 返回:如果任务无法取消,通常是因为它已经完成则false否则true。如果两个或多个线程导致任务被取消,则其中至少有一个线程返回true。实现可能会提供更强的保证。
**/
boolean cancel(boolean mayInterruptIfRunning);
/**
* 如果此任务在正常完成之前被取消,则返回 true 此任务。
* 返回:true 如果此任务在完成之前被取消
**/
boolean isCancelled();
/**
* 如果此任务已完成,则返回 true 。完成可能是由于正常终止、异常或取消 —— 在所有这些情况下,此方法都将返回 true.
* 返回:true 如果此任务已完成
**/
boolean isDone();
/**
* 如有必要,请等待计算完成,然后检索其结果。
* 返回:计算结果
* 抛出:CancellationException – 如果计算被取消,ExecutionException – 如果计算引发异常,InterruptedException – 如果当前线程在等待时被中断
**/
V get() throws InterruptedException, ExecutionException;
/**
* 如有必要,最多等待给定时间以完成计算,然后检索其结果(如果可用)。
* 参数:timeout – 等待的最长时间 unit – timeout 参数的时间单位
* 返回:计算结果
* 抛出:
* CancellationException – 如果计算被取消
* ExecutionException – 如果计算引发异常
* InterruptedException – 如果当前线程在等待时被中断
* TimeoutException – 如果等待超时
**/
V get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException;默认方法(jdk19及之后才有)如下:
/**
* 返回计算结果,无需等待。
* 此方法适用于调用方知道任务已成功完成的情况,例如,在筛选 Future 对象流以获取成功任务并使用 Map作获取结果流时。
* results = futures.stream()
* .filter(f -> f.state() == Future.State.SUCCESS)
* .map(Future::resultNow)
* .toList();
* 返回:计算结果
* 抛出:IllegalStateException – 如果任务尚未完成或任务未完成并显示结果
* 实现要求:默认实现调用 isDone() 以测试任务是否已完成。如果完成,它将调用 get() 以获取结果。
**/
default V resultNow() {
if (!isDone())
throw new IllegalStateException("Task has not completed");
boolean interrupted = false;
try {
while (true) {
try {
return get();
} catch (InterruptedException e) {
interrupted = true;
} catch (ExecutionException e) {
throw new IllegalStateException("Task completed with exception");
} catch (CancellationException e) {
throw new IllegalStateException("Task was cancelled");
}
}
} finally {
if (interrupted) Thread.currentThread().interrupt();
}
}工作原理图:线程池处理任务
flowchart TD
A["提交任务 execute/submit"] --> B{"核心线程是否空闲或可创建"}
B -- "是" --> C["核心线程执行任务"]
B -- "否" --> D{"任务队列是否已满"}
D -- "否" --> E["任务进入阻塞队列"]
E --> F["工作线程从队列取任务"]
D -- "是" --> G{"是否还能创建非核心线程"}
G -- "是" --> H["创建临时线程执行"]
G -- "否" --> I["执行拒绝策略"]
C --> J["任务完成"]
F --> J
H --> J零基础理解线程池时,重点看三个容量:核心线程数、队列长度、最大线程数。任务不是无限创建线程执行,而是按规则在线程、队列、拒绝策略之间流转。
为什么是“核心线程 -> 队列 -> 最大线程 -> 拒绝策略”
很多初学者会疑惑:既然最大线程数更大,为什么不一开始就把线程创建到最大?
线程池这样设计是为了在吞吐、资源和稳定性之间折中:
| 阶段 | 设计目的 | 如果没有这一层会怎样 |
|---|---|---|
| 核心线程 | 维持一批长期工作线程,避免频繁创建销毁 | 每个任务都临时建线程,成本高 |
| 阻塞队列 | 平滑短时间流量峰值,让任务排队等待 | 流量稍微抖动就创建大量线程 |
| 最大线程 | 队列也顶不住时临时扩容处理压力 | 峰值时只能排队,延迟快速升高 |
| 拒绝策略 | 系统已经没有处理能力时明确失败 | 无限堆积直到 OOM 或拖垮服务 |
流程可以这样理解:
flowchart TD
A["任务到来"] --> B{"核心线程还有名额"}
B -- "有" --> C["创建或复用核心线程"]
B -- "没有" --> D{"队列还能放"}
D -- "能" --> E["进入队列等待"]
D -- "不能" --> F{"还能创建非核心线程"}
F -- "能" --> G["创建非核心线程处理峰值"]
F -- "不能" --> H["拒绝策略保护系统"]为什么不直接创建到最大线程数:线程不是越多越好。线程太多会带来上下文切换、内存栈空间占用、数据库连接竞争、下游接口放大压力。线程池的意义不是“尽量多跑”,而是“在系统能承受的范围内稳定处理”。
生产建议:显式创建 ThreadPoolExecutor
生产环境不建议无脑使用 Executors.newFixedThreadPool()、newCachedThreadPool(),因为它们隐藏了队列或线程数量风险。更推荐显式声明线程数、队列和拒绝策略。
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class ThreadPoolExecutorDemo {
public static void main(String[] args) {
ThreadPoolExecutor executor = new ThreadPoolExecutor(
4,
8,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(200),
runnable -> {
Thread thread = new Thread(runnable);
thread.setName("collect-worker-" + thread.threadId());
return thread;
},
new ThreadPoolExecutor.CallerRunsPolicy()
);
executor.execute(() -> System.out.println("采集任务执行"));
executor.shutdown();
}
}这个 Demo 的关键点:
corePoolSize=4:平时维持 4 个工作线程。maximumPoolSize=8:峰值时最多扩到 8 个线程。ArrayBlockingQueue<>(200):队列有明确边界,避免无限堆积。CallerRunsPolicy:线程池满时让提交任务的线程自己执行,形成反压。
如果采集任务依赖数据库连接池,线程池最大线程数不能大于下游可承受能力太多。否则线程多了也只是一起等连接、等接口,吞吐未必提升,延迟和故障传播反而更严重。
线程池是不是每个线程都有一个任务队列
不是。ThreadPoolExecutor 通常是“一个线程池共享一个阻塞队列”,不是每个工作线程一个队列。
flowchart TD
A["ThreadPoolExecutor"] --> B["共享 BlockingQueue"]
A --> C["Worker-1"]
A --> D["Worker-2"]
A --> E["Worker-3"]
B --> C
B --> D
B --> E任务先提交到线程池,线程池根据核心线程、队列、最大线程和拒绝策略安排执行。工作线程执行完当前任务后,会继续从共享队列中取任务。
如果误以为“每个线程都有自己的队列”,就会错误理解任务调度:线程池不是给某个线程固定分配一串任务,而是一组 Worker 共同消费队列。
核心线程和非核心线程有没有永久身份标识
面试里经常问:线程池会不会给线程打标记,告诉它“你是核心线程”“你是非核心线程”?
更准确的理解是:线程池主要维护的是当前工作线程数量、核心线程数、最大线程数、空闲超时时间等规则,而不是给每个 Worker 永久贴一个“核心/非核心”的身份标签。
当线程数量超过 corePoolSize 时,空闲超过 keepAliveTime 的 Worker 可能被回收。是否回收取决于当前线程池数量和是否允许核心线程超时。
flowchart TD
A["Worker 空闲等待任务"] --> B{"是否超时"}
B -- "否" --> C["继续等待"]
B -- "是" --> D{"当前线程数是否大于 corePoolSize"}
D -- "是" --> E["回收该 Worker"]
D -- "否" --> F{"allowCoreThreadTimeOut 是否开启"}
F -- "是" --> E
F -- "否" --> C结论:所谓核心线程,更像是线程池保留线程数量的规则;所谓非核心线程,是超过核心数量后为了处理峰值临时增加的 Worker,空闲超时后会被回收。
代码 Demo:固定线程池提交任务
ExecutorService executorService = Executors.newFixedThreadPool(3);
Future<Integer> future = executorService.submit(() -> {
Thread.sleep(1000);
return 1 + 2;
});
System.out.println(future.get());
executorService.shutdown();submit 会返回 Future,适合需要结果的任务;execute 只提交任务,不直接返回执行结果。
常见风险
| 风险 | 后果 | 建议 |
|---|---|---|
| 使用无界队列 | 任务堆积导致内存持续上涨 | 估算容量,使用有界队列 |
不处理 Future.get() 异常 | 任务失败被吞掉,业务以为成功 | 捕获 ExecutionException 并记录上下文 |
忘记 shutdown() | 应用无法正常退出,线程泄漏 | 应用停止时关闭线程池 |
| 在线程池任务里阻塞等待同一个线程池任务 | 可能出现线程池饥饿死锁 | 拆分线程池或避免同步等待 |
盲目使用 Executors.newCachedThreadPool() | 线程数量可能失控 | 生产环境优先显式创建 ThreadPoolExecutor |
