线程池、队列、拒绝与隔离
ThreadPoolExecutor 把提交的任务分配给 worker,暂时无法执行的任务进入队列。当队列不能接收、worker 也无法增加时,拒绝策略决定提交者得到什么结果。
已提交任务
├─ 执行中:受 worker 数量和下游资源约束
├─ 排队中:占用等待时间与堆内存
└─ 未接受:抛异常、调用者执行或按策略丢弃线程数、队列容量和拒绝方式需要一起配置。只增大 maximumPoolSize,而使用几乎总能接收任务的无界队列,通常不会触发预期的扩线程分支。
用有界任务复现饱和
完整的 1 核心、2 最大、1 队列槽实验
下面的 Saturation 用闩锁暂时挡住任务。前三次提交分别创建核心 worker、入队、创建第二个 worker,第四次提交触发拒绝。finally 释放闩锁并关闭池:
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
public final class Saturation {
public static void main(String[] args) throws InterruptedException {
CountDownLatch release = new CountDownLatch(1);
AtomicInteger names = new AtomicInteger();
AtomicReference<Throwable> failure = new AtomicReference<>();
ThreadPoolExecutor pool = new ThreadPoolExecutor(
1, 2, 10, TimeUnit.SECONDS, new ArrayBlockingQueue<>(1),
task -> new Thread(task, "orders-" + names.incrementAndGet()),
new ThreadPoolExecutor.AbortPolicy());
try {
Runnable task = () -> {
try {
if (!release.await(10, TimeUnit.SECONDS)) throw new AssertionError("release timeout");
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
catch (Throwable e) { failure.compareAndSet(null, e); }
};
pool.execute(task);
pool.execute(task);
pool.execute(task);
boolean rejected = false;
try { pool.execute(task); }
catch (RejectedExecutionException expected) { rejected = true; }
if (!rejected || pool.getPoolSize() != 2 || pool.getQueue().size() != 1)
throw new AssertionError("saturation");
System.out.println("pool=2 queue=1 rejected=true");
} finally {
release.countDown();
pool.shutdown();
if (!pool.awaitTermination(5, TimeUnit.SECONDS)) {
pool.shutdownNow();
if (!pool.awaitTermination(5, TimeUnit.SECONDS)) throw new AssertionError("termination");
}
}
if (failure.get() != null) throw new AssertionError(failure.get());
}
}Linux 普通用户在可写目录保存为 Saturation.java,使用完整 JDK,安装见Java 版本基线。设置实际 JAVA_HOME;代码以 Java 17 为编译目标,在 Temurin 25.0.4+7 上运行:
export JAVA_HOME=/opt/jdk-25
unset JAVA_TOOL_OPTIONS JDK_JAVA_OPTIONS _JAVA_OPTIONS
LAB_OUT=$(mktemp -d /tmp/pool-learning.XXXXXX)
"$JAVA_HOME/bin/javac" --release 17 -Xlint:all -Werror -d "$LAB_OUT" Saturation.java
"$JAVA_HOME/bin/java" -cp "$LAB_OUT" Saturation预期 pool=2 queue=1 rejected=true,随后进程自行结束。闩锁使前三个任务不会在第四次提交前正常完成,避免靠“慢任务大概还没结束”判断饱和。时间限制只是实验失败保护,不是生产线程参数。
若输出与预期不同,检查是否改动了队列、线程数或释放顺序。出现 release timeout 表示实验环境或程序执行严重延迟,应判失败并检查完整输出,而不是只保留饱和一行。
执行接口与配置对象
Executor 只定义 execute,ExecutorService 增加 submit、Future 和生命周期管理,ThreadPoolExecutor 提供可配置的线程与队列实现。Runnable 不直接返回结果,Callable 可返回值或抛异常;submit 将它们包装成可以观察完成状态的 Future。
| 参数 | 作用 |
|---|---|
| corePoolSize | 新任务提交时优先创建 worker 的目标数量 |
| maximumPoolSize | 入队不能接收时允许的 worker 上限 |
| keepAliveTime / unit | 可回收空闲 worker 的等待期限 |
| workQueue | 尚未开始任务的存放与排序方式 |
| threadFactory | 线程命名、创建与异常处理 |
| handler | 未能接受任务时的处理 |
核心 worker 默认按需创建,可通过 prestartCoreThread / prestartAllCoreThreads 提前启动。允许核心线程超时需要启用 allowCoreThreadTimeOut 且 keepAliveTime 为正。配置与监控契约见 ThreadPoolExecutor API。
这些参数允许部分动态调整,但调整不等于立即中止已运行任务。减小最大线程数后,超额线程按实现退出条件回收;已有队列实例也不会随一个“容量配置值”自动换成新队列。
execute 怎样决定执行、入队和拒绝
创建失败还要进入后续分支
execute(command)
├─ worker 数小于 core
│ └─ addWorker(command, true) 成功 → 返回
│ 失败 → 继续判断
├─ 池为 RUNNING 且 offer 成功
│ ├─ 复查已关闭且 remove 成功 → 拒绝
│ └─ 否则若 worker 数为 0 → 尝试补 worker 消费队列
└─ 无法入队
├─ addWorker(command, false) 成功 → 执行首任务
└─ 失败 → 拒绝固定 OpenJDK jdk-25+36 ThreadPoolExecutor 实现用这个顺序处理并发提交。addWorker 可能因状态变化、数量竞争或线程创建问题失败,因此“低于核心数”不保证第一次创建一定完成。
入队后的 recheck 处理关闭竞态:提交者刚 offer,另一线程可能已 shutdown。若能移除该任务就拒绝;若任务已经被 worker 取走,不能再承诺拒绝成功。检测到没有 worker 时尝试补建,使队列得到消费机会。
ThreadFactory 返回 null 会妨碍创建 worker,未必让每次提交都抛异常;任务可能进入队列却无人消费。工厂应稳定创建具有可辨识名称的线程,创建失败必须报告。不要在工厂里执行慢网络初始化。
ctl 与 Worker 的不同职责
ctl 将运行状态和 worker 计数放进一个原子整数,便于同时校验数量与生命周期。应用通过公共指标与 isShutdown/isTerminated 观察,不读取私有位编码。
| 状态 | 新任务 | 队列与正在运行的任务 |
|---|---|---|
| RUNNING | 接受 | 正常处理 |
| SHUTDOWN | 拒绝 | 排空此前接受的队列,运行任务继续 |
| STOP | 拒绝 | 不再启动队列任务,尝试中断正在运行者 |
| TIDYING | 不接受 | worker 已归零,执行终止钩子 |
| TERMINATED | 不接受 | 终止钩子完成 |
Worker 持有 thread、firstTask 和完成计数,还用自己的同步状态区分忙碌与空闲。shutdown 可以只中断空闲 worker 使其从取任务等待中醒来;运行中的任务不会因此直接被强制终止。shutdownNow 则会尝试中断全部工作线程。
runWorker 执行首任务后循环 getTask。getTask 根据池状态、worker 数与超时配置选择 take 或带时限 poll,决定空闲回收。栈在 getTask/BlockingQueue.take 常是正常空闲;栈在数据库、HTTP 客户端或业务锁中,才说明 worker 正在等待任务依赖。
beforeExecute、afterExecute 和 terminated 可用于轻量观测。钩子抛异常会影响 worker,采集指标、清理上下文都应避免阻塞和不受控异常。重建 worker 可以恢复执行资源,但不会自动重试刚刚失败的业务任务。
任务的结果、异常与取消
下载线程池实验包,解压进入 threadpool-queue-rejection,沿用 LAB_OUT:
"$JAVA_HOME/bin/javac" --release 17 -Xlint:all -Werror -d "$LAB_OUT" *.java
"$JAVA_HOME/bin/java" -cp "$LAB_OUT" PoolLab异常部分输出:
execute=uncaught submit=ExecutionException replacement-worker=trueexecute 直接执行的 Runnable 抛异常时,异常可以到达线程的未捕获异常处理器,当前 worker 异常退出;实验随后提交任务,验证池继续工作。submit 通常使用 FutureTask 保存异常,Future.get 才以 ExecutionException 交还原因。不能依赖未捕获异常处理器发现所有 submit 失败。
afterExecute 接收的 Throwable 也可能为 null,因为 FutureTask 已将异常存入结果。自定义观测若要解包,先判断任务是 Future 且已完成,再获取结果并分别记录取消、执行失败和中断,不能在钩子中无条件等待尚未完成的 Future。
Future.get(timeout) 限制调用者等待结果的时间,不会自动取消后台任务。Future.cancel(true) 对可中断的执行器任务发出取消请求,任务仍需响应;已经发生的数据库写入也不会被自动撤销。取消语义见 Future API。
提交前后的内存交接与 Future.get 的可见性由 ExecutorService 契约定义。线程池复用线程,不意味着任务可以继续使用上一个请求遗留的 ThreadLocal;上下文传播与恢复见异步上下文、Future 与虚拟线程。
排队容量与拒绝结果
队列先决定扩容机会
| 队列 | 提交压力怎样传递 |
|---|---|
| ArrayBlockingQueue / 显式容量 LinkedBlockingQueue | 满后尝试增加 worker,再拒绝 |
| 默认 LinkedBlockingQueue | 大多数任务继续排队,maximumPoolSize 很少触发 |
| SynchronousQueue | 无存储槽,不能直接交接时尝试增加 worker |
| PriorityBlockingQueue | 优先出队但逻辑无界,低优先级可能长期等待 |
Executors.newFixedThreadPool 与 newSingleThreadExecutor 使用无界队列;newCachedThreadPool 可快速扩展线程。工厂方便不代表默认容量适用于所有服务,具体定义见 Executors API。
优先队列还要求任务具有可比较关系或提供比较器。submit 包装后的 FutureTask 不会自动继承业务 Runnable 的优先级接口;需要明确的任务包装与执行器扩展,不能只替换构造器里的队列。
队列容量表示允许多少尚未开始的任务。它应同时服从等待预算与堆预算:一个任务可能捕获请求体、回调与 Future,平均体积和长尾体积都影响内存。队列计数之外还要记录最老任务年龄与开始前等待时间。相关队列操作见并发容器。
四种策略与未完成 Future
| 策略 | 运行中饱和 | 关闭后 | 结果风险 |
|---|---|---|---|
| AbortPolicy | 抛 RejectedExecutionException | 抛异常 | 调用方可以明确处理未接受 |
| CallerRunsPolicy | 提交线程执行 | 丢弃 | 入口线程被占用,关闭时没有执行结果 |
| DiscardPolicy | 静默丢弃 | 静默丢弃 | submit 返回的 Future 可能一直未完成 |
| DiscardOldestPolicy | 丢队列头后重试 | 丢弃 | 旧任务回执失去处理者,优先队列头未必最旧 |
CallerRunsPolicy 的关闭行为见官方 API。它只把工作转移给当前提交者:HTTP 线程会变慢,持锁提交者可能在锁内执行整个任务,事件循环也可能被阻塞。若原任务又提交依赖任务,调用栈和等待关系还会改变。
PoolLab 在关闭池上使用 CallerRunsPolicy 和 DiscardPolicy 提交任务,断言返回的 Future 尚未完成,然后显式取消以收尾:
closed-caller-runs future-incomplete=true explicitly-cancelled=true
discard future-incomplete=true explicitly-cancelled=true因此,需要结果的任务不能配一个无回执的丢弃策略。自定义拒绝处理应快速计数并抛出可识别异常,或由统一提交封装完成失败 Future;避免在提交线程同步上报网络、无限重试。必须持久保留的任务应先进入可靠存储,再由执行器处理,进程内队列不提供宕机恢复。
线程数与下游并发共同设限
CPU 密集工作可以从容器实际可用处理器数附近开始压测;阻塞工作需结合等待时间增加并发,但上限还受连接池、下游配额、单任务内存和文件句柄限制。增加线程只会把等待从执行队列搬到连接池时,吞吐不会相应增长。
常用估算为“可用 CPU × 目标利用率 × (1 + 等待/计算时间)”。它假设工作特征相对稳定,只用于选择压测起点。若某类请求同时调用多个下游或重试多次,必须计入扇出与重试放大。
任务开始前可以根据单调时钟预算拒绝已过期工作。若包装 FutureTask,跳过 run 时仍需将对应 Future 完成失败或取消,否则调用者会永久等待。暂停低价值生产者、限制入口并发,比让陈旧任务填满队列更直接。
定时执行与依赖隔离
ScheduledThreadPoolExecutor 使用延迟队列,主要由 corePoolSize 决定并发能力,maximumPoolSize 对这种无界延迟队列没有通常期待的扩容效果。schedule 执行一次;scheduleAtFixedRate 以计划频率推进,scheduleWithFixedDelay 从上一次结束后计算间隔。同一周期任务不会重叠执行,执行过慢会使后续开始延后。
周期任务若向外抛异常,后续执行被抑制。PoolLab 的周期任务第一次便抛异常,并从其 Future 读取失败:
periodic-runs=1 failure-stops-series=true行为见 ScheduledExecutorService API。需要持续调度的业务应在任务内部识别可恢复异常、记录失败并结束本轮;不可恢复错误应触发停用和告警,不无限吞掉所有 Throwable。
取消的延迟任务默认可能保留到延迟到期,可使用 setRemoveOnCancelPolicy(true) 及时移除。关闭后是否继续已有延迟/周期任务由对应策略控制,详见 ScheduledThreadPoolExecutor API。长延迟任务没有必要捕获整个请求上下文,进程内调度也无法跨重启保留计划。
所有 worker 被父任务占满,而父任务向同池提交子任务并 get,会形成依赖环:
父任务占用全部 worker → 子任务排队
↑ │
└── 等子任务完成 ────┘解决需要打断同步依赖:直接执行适当的小子任务、改成非阻塞完成组合,或在有独立容量依据时分离执行资源。仅加线程会推迟再次饥饿的负载阈值。稳定复现与线程采集见并发诊断。
隔离可以按慢下游、CPU 工作与后台批处理划分,不必每个方法创建一个池。组件拥有执行器并统一关闭;总线程、总排队内存与共享数据库连接仍需一起预算。
关闭流程与故障处理
shutdownNow 返回任务仍要处理回执
先停止新的提交入口,再 shutdown 等待已接受任务结束。超过业务允许的排空时间后,shutdownNow 尝试中断运行任务,并返回尚未开始的 Runnable。它不保证线程立即退出,也不保证返回队列中的 Future 已被取消。
PoolLab 占住一个 worker,排入一个 FutureTask 后 shutdownNow;取出任务仍未完成,实验逐项取消返回的 Future,再验证运行任务观察到中断:
shutdownNow-pending=1 queued-future-cancelled=true running-interrupted=true应用可以据返回列表记录任务 ID、完成失败回执或转交可靠补偿。若任务包装器隐藏内部 Future,还需由包装协议暴露取消动作。不要将 shutdownNow 的调用返回当作所有工作已经结束,应继续 awaitTermination 并检查实际任务去向。
Java 19 起 ExecutorService 支持 AutoCloseable,close 等待终止,不能当成带短时限的关闭操作。Java 17 示例使用显式关闭,避免版本差异,也方便设置排空预算。
观察等待、执行和结果
getPoolSize、getActiveCount、getQueue().size 和 getCompletedTaskCount 是运行估计。完成任务计数包含执行结束的任务,业务成功要从结果通道独立记录。
提交封装可以保存 submittedAt,任务开始记录 startedAt,finally 记录 finishedAt,以计算排队与执行耗时。另行记录拒绝、失败和取消,避免把每个计数都相加为“总任务”:例如拒绝发生在接受之前,取消又可能发生在排队或运行阶段。
| 现象 | 首要判断 | 处理后检查 |
|---|---|---|
| queue 增长,worker 一直是 core | 队列是否无界 | 有界配置下实际扩展和拒绝可见 |
| active 满、CPU 低 | 下游等待、锁、同池子任务 | 打断依赖或恢复下游后队列年龄下降 |
| 关闭阶段突然拒绝 | 提交入口是否仍活跃 | 先停入口,再排空与关闭 |
| completed 增长却缺业务结果 | Future 失败未消费、静默丢弃 | 每个接受任务最终有结果或明确取消 |
| queue 非空、worker 为 0 | 工厂失败、资源耗尽与池状态 | 线程恢复后实际消费,不只修改 max |
| 周期任务突然不再运行 | Future 异常或已取消 | 修复异常后按明确策略重新调度 |
采集活跃线程、CPU 与池饥饿实验见并发诊断。限流、暂停后台生产者与缩短无价值排队可先减轻负载;扩线程前确认下游仍有可用处理能力。
完整运行与清理
在解压目录设置 JAVA_HOME 后运行 bash run.sh,严格编译并执行饱和、异常、静默拒绝、关闭和定时断言,最后输出 PASS pool-lab。Docker 运行由有权限的 Linux 主机用户执行:
LAB_SOURCE="$(pwd -P)"
JDK_IMAGE=eclipse-temurin:25.0.4_7-jdk
docker pull "$JDK_IMAGE"
docker run --pull never --rm --network none --read-only \
--user 10001:10001 --cap-drop ALL --security-opt no-new-privileges \
--mount "type=bind,src=$LAB_SOURCE,dst=/lab,readonly" \
--tmpfs /tmp:rw,nosuid,nodev,size=128m,mode=1777 \
"$JDK_IMAGE" bash /lab/run.shUID 10001 只读取源码,产物写 tmpfs。离线环境从可信渠道准备并校验镜像归档后导入;Java 17 对照镜像是 eclipse-temurin:17.0.20_8-jdk。脚本每个 Java 进程外部限制 20 秒,超时应调查线程和资源,不作为成功结束。
确认实验进程结束后,回收手工目录:
case "$LAB_OUT" in
/tmp/pool-learning.*) rm -r -- "$LAB_OUT" ;;
*) printf '%s\n' '保留未知目录' ;;
esac
unset LAB_OUT权威资料与规范地址
按方法语义查 API,按实现字段查固定版本源码。
| 资料 | 完整地址 |
|---|---|
| ThreadPoolExecutor API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ThreadPoolExecutor.html |
| OpenJDK jdk-25+36 ThreadPoolExecutor 实现 | https://github.com/openjdk/jdk/blob/jdk-25%2B36/src/java.base/share/classes/java/util/concurrent/ThreadPoolExecutor.java |
| Future API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/Future.html |
| ExecutorService 契约 | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ExecutorService.html |
| Executors API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/Executors.html |
| 官方 API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ThreadPoolExecutor.CallerRunsPolicy.html |
| ScheduledExecutorService API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ScheduledExecutorService.html |
| ScheduledThreadPoolExecutor API | https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/ScheduledThreadPoolExecutor.html |
