Skip to content

ExecutorService

本章以企业项目最常见的 JDK 7/JDK 8 线程池模型 为主线讲解:ExecutorExecutorServiceThreadPoolExecutorCallableFuture。文中涉及 JDK 19/21 之后新增的 Future 默认方法时,只作为版本演进补充,不作为基础学习前提。

ExecutorService 是线程池体系里的核心执行器接口,继承自 Executor,常见实现是 ThreadPoolExecutor。它提供了两类常用提交方式:不关心返回值时用 execute(),需要结果、异常或取消时用 submit() 返回 Future

为什么需要 ExecutorService

如果每来一个任务都 new Thread(),短期看起来简单,但高并发时会出现三个问题:线程创建销毁成本高、线程数量不可控、任务失败和排队策略不可管理。ExecutorService 把“任务提交”和“线程执行”分开,让线程可以复用,也让队列、拒绝策略、关闭流程变成可控配置。

mermaid
flowchart TD
    A["业务提交任务"] --> B["ExecutorService"]
    B --> C{"是否有可用线程"}
    C -- "有" --> D["复用工作线程执行"]
    C -- "没有" --> E["进入任务队列或触发拒绝策略"]
    D --> F["任务完成后线程继续复用"]

如果不会线程池原理,最容易把线程池当成“无限异步工具”。结果是队列堆满、内存升高、请求延迟变长,最后触发拒绝策略或服务整体不可用。

execute方法

execute() 方法是Executor接口定义的方法,当你只需要提交一个简单的Runnable任务不关心其返回值时可以使用这个方法。定义如下:

java
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 对象可以用来检查任务是否完成、等待任务完成并获取其结果。抽象方法如下:

java
/** 提交一个返回值的任务以供执行,并返回一个 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对象,其常见抽象方法如下:

java
/**
 * 尝试取消此任务的执行。如果任务已经完成或取消,或者由于其他原因无法取消,则此方法无效。否则,如果此任务在调用时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及之后才有)如下:

java
/**
 * 返回计算结果,无需等待。
 * 此方法适用于调用方知道任务已成功完成的情况,例如,在筛选 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();
        }
    }

工作原理图:线程池处理任务

mermaid
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 或拖垮服务

流程可以这样理解:

mermaid
flowchart TD
    A["任务到来"] --> B{"核心线程还有名额"}
    B -- "有" --> C["创建或复用核心线程"]
    B -- "没有" --> D{"队列还能放"}
    D -- "能" --> E["进入队列等待"]
    D -- "不能" --> F{"还能创建非核心线程"}
    F -- "能" --> G["创建非核心线程处理峰值"]
    F -- "不能" --> H["拒绝策略保护系统"]

为什么不直接创建到最大线程数:线程不是越多越好。线程太多会带来上下文切换、内存栈空间占用、数据库连接竞争、下游接口放大压力。线程池的意义不是“尽量多跑”,而是“在系统能承受的范围内稳定处理”。

生产建议:显式创建 ThreadPoolExecutor

生产环境不建议无脑使用 Executors.newFixedThreadPool()newCachedThreadPool(),因为它们隐藏了队列或线程数量风险。更推荐显式声明线程数、队列和拒绝策略。

java
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 的关键点:

  1. corePoolSize=4:平时维持 4 个工作线程。
  2. maximumPoolSize=8:峰值时最多扩到 8 个线程。
  3. ArrayBlockingQueue<>(200):队列有明确边界,避免无限堆积。
  4. CallerRunsPolicy:线程池满时让提交任务的线程自己执行,形成反压。

如果采集任务依赖数据库连接池,线程池最大线程数不能大于下游可承受能力太多。否则线程多了也只是一起等连接、等接口,吞吐未必提升,延迟和故障传播反而更严重。

线程池是不是每个线程都有一个任务队列

不是。ThreadPoolExecutor 通常是“一个线程池共享一个阻塞队列”,不是每个工作线程一个队列。

mermaid
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 可能被回收。是否回收取决于当前线程池数量和是否允许核心线程超时。

mermaid
flowchart TD
    A["Worker 空闲等待任务"] --> B{"是否超时"}
    B -- "否" --> C["继续等待"]
    B -- "是" --> D{"当前线程数是否大于 corePoolSize"}
    D -- "是" --> E["回收该 Worker"]
    D -- "否" --> F{"allowCoreThreadTimeOut 是否开启"}
    F -- "是" --> E
    F -- "否" --> C

结论:所谓核心线程,更像是线程池保留线程数量的规则;所谓非核心线程,是超过核心数量后为了处理峰值临时增加的 Worker,空闲超时后会被回收。

代码 Demo:固定线程池提交任务

java
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