Java 基础体系 · 第 13/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java 线程与 Executor:生命周期、线程池、队列、Future 和取消

Java 并发代码通常同时包含两类对象:

  • 任务:要执行的逻辑,例如实现了 RunnableCallable<T> 的对象;
  • 执行者:负责安排任务在哪个线程上执行,例如 ThreadExecutorExecutorService

Thread 解决的是“一个线程如何创建和运行”;Executor 解决的是“任务如何交给执行机制,以及如何管理执行资源”。线程池、工作队列、Future 和取消,都是围绕这两个层次展开的。

本文以 Java 25 为范围,区分 Java 规范保证、标准 API 契约和常见实现行为。


1. 从任务到线程:先区分三个生命周期

一个任务至少经历三个不同生命周期:

  1. 任务对象生命周期:对象被创建、提交、运行、结束;
  2. 线程生命周期:线程从 NEW 进入可运行状态,最后变为 TERMINATED
  3. 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();

线程中的未捕获异常会终止该线程,并交给未捕获异常处理器;它不会像普通方法返回值一样自动传播到创建线程的调用栈。

在线程池中需要额外区分 executesubmit,后文会说明:submit 通常把异常保存到 Future 中,而不是立即交给线程的未捕获异常处理器。

3.3 中断是协作式取消信号

Thread.interrupt() 的语义取决于目标线程当前做什么:

  1. 如果目标线程正在 sleepwaitjoin 等可中断阻塞中,通常会收到 InterruptedException,并清除中断标志;
  2. 如果目标线程没有处于这些阻塞调用中,中断标志会被设置;
  3. 如果目标线程完全忽略中断,线程不会因此自动停止。
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 增加生命周期和结果管理

ExecutorServiceExecutor 之上增加了:

  • submit:提交 RunnableCallable<T>
  • Future:观察任务结果、异常和取消状态;
  • shutdown:停止接收新任务;
  • shutdownNow:尝试中断执行中的任务,并返回尚未开始的任务;
  • awaitTermination:等待 Executor 终止;
  • invokeAllinvokeAny:批量执行任务。

典型生命周期:

stateDiagram-v2
    [*] --> RUNNING
    RUNNING --> SHUTDOWN: shutdown()
    RUNNING --> STOP: shutdownNow()
    SHUTDOWN --> TIDYING: 已提交任务全部完成
    STOP --> TIDYING: 工作线程结束
    TIDYING --> TERMINATED: terminated()

ExecutorServiceisShutdown()isTerminated() 含义不同:

  • isShutdown():已经不再接受新任务;
  • isTerminated():关闭流程已完成,工作线程已退出。

调用 shutdown() 后,isShutdown() 可能立刻为 true,但 isTerminated() 仍然为 false


6. 线程池的核心组成

典型平台线程池包含:

  1. 工作线程:真正执行任务;
  2. 工作队列:保存暂时没有线程执行的任务;
  3. 线程工厂:创建和命名线程;
  4. 拒绝策略:线程和队列都无法接收任务时的处理方式;
  5. 生命周期控制:关闭、等待和中断。

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 请求线程、事件循环线程或关键调度线程,可能把耗时任务带入不应阻塞的线程。

DiscardPolicyDiscardOldestPolicy 只有在业务明确允许丢失任务时才可使用。否则,静默丢弃会让故障变成数据缺失。


9. executesubmitFuture

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();
        }
    }
}

执行过程通常是:

  1. 单线程池创建一个工作线程;
  2. 任务循环运行;
  3. 主线程调用 cancel(true)
  4. 工作线程在 sleep 中收到 InterruptedException
  5. 任务恢复中断状态并退出;
  6. finally 执行清理;
  7. Executor 关闭。

11.2 不响应取消的任务

executor.submit(() -> {
    while (true) {
        calculate();
    }
});

如果 calculate() 不检查中断,也不调用可中断阻塞操作,那么:

future.cancel(true);

通常只能使 Future 进入取消状态,而任务仍可能继续消耗 CPU。

这也是 shutdownNow() 不能保证立即停止所有任务的根本原因:它的主要机制也是中断工作线程,而中断是协作式的。


12. shutdownshutdownNow 和资源释放

12.1 正常关闭

executor.shutdown();

try {
    if (!executor.awaitTermination(30, TimeUnit.SECONDS)) {
        executor.shutdownNow();
    }
} catch (InterruptedException e) {
    executor.shutdownNow();
    Thread.currentThread().interrupt();
}

这个关闭流程的含义是:

  1. shutdown():不再接受新任务;
  2. 已提交任务继续执行;
  3. awaitTermination 等待工作线程结束;
  4. 超时后调用 shutdownNow() 尝试中断剩余任务;
  5. 如果等待线程自身被中断,也应继续关闭并恢复中断状态。

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. 批量任务:invokeAllinvokeAny

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

关闭是单向生命周期变化。应用停止时通常应:

  1. 停止产生新任务;
  2. 调用 shutdown()
  3. 等待有限时间;
  4. 必要时调用 shutdownNow()
  5. 记录仍未结束的任务;
  6. 恢复被中断的关闭线程状态。

如果任务会向另一个 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、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。