线程池、队列、拒绝与隔离:过载怎样被控制在边界内
接口 P99 已经超过 3 秒,线程池却没有拒绝,监控甚至显示“最大线程数 200,当前线程只有 20”。这不是容量充足,而可能是无界队列把任务全部收下了:maximumPoolSize 根本没有机会生效,用户早已超时的任务仍在队列里等待,内存和下游压力继续累积。
线程池不是“把任务异步一下”的工具,而是一条容量边界。它决定多少任务同时执行,多少任务可以等待,等待多久后失去价值,满了以后谁承担压力,服务关闭时哪些任务能完成。把 ThreadPoolExecutor 的运行链读懂后,参数才能从业务预算和下游容量推出来。
先用 1 个核心线程、1 个队列槽复现饱和
ThreadPoolExecutor pool = new ThreadPoolExecutor(
1,
2,
30,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1),
namedThreadFactory("order-query"),
new ThreadPoolExecutor.AbortPolicy()
);
for (int i = 1; i <= 4; i++) {
int taskId = i;
try {
pool.execute(() -> slowTask(taskId));
} catch (RejectedExecutionException e) {
System.out.println("rejected=" + taskId);
}
}预期链路是:任务 1 创建核心 worker;任务 2 入队;任务 3 在队列满后创建非核心 worker;任务 4 在线程和队列都满时被拒绝。这个实验比背七个构造参数更有用,因为它把提交顺序画出来了。
入队后要复查,是因为提交期间线程池可能从 RUNNING 进入 SHUTDOWN;如果不能继续接收,已入队任务要移除并拒绝。如果池仍运行但没有 worker,还要补一个 worker 来消费队列。并发源码里的二次检查往往就是在封这类状态窗口。
ctl 把生命周期和 worker 数绑在一个原子状态里
OpenJDK ThreadPoolExecutor 使用原子整数 ctl 编码运行状态与 worker 数。高位描述 RUNNING、SHUTDOWN、STOP、TIDYING、TERMINATED 等生命周期,低位描述当前 worker 数。理解它是为了反推现象,不是要求业务读取私有字段。
| 状态 | 新任务 | 队列任务 | 正在运行任务 |
|---|---|---|---|
| RUNNING | 接收 | 处理 | 继续 |
| SHUTDOWN | 拒绝 | 处理完 | 继续 |
| STOP | 拒绝 | 不再处理 | 请求中断 |
| TIDYING | 无任务和 worker,执行收尾 | 无 | 无 |
| TERMINATED | 已完成 terminated | 无 | 无 |
worker 数量和运行状态必须一起看。poolSize=0 可能是空闲线程已回收,也可能是工厂创建线程失败;队列里有任务但没有 worker 是异常现场;进程发布中出现拒绝,还要判断池是否已进入关闭状态,而不是只调大最大线程。
Worker 为什么既是任务执行者也是锁
worker 持有线程,并通过 runWorker 循环执行首个任务、再从 getTask() 取后续任务。getTask 根据核心线程、最大线程、超时回收和池状态决定用 take 还是 poll,以及 worker 是否该退出。
execute
-> addWorker
-> Worker(Thread + firstTask)
-> Thread.start
-> runWorker
-> beforeExecute
-> task.run
-> afterExecute
-> getTask线程栈停在 ThreadPoolExecutor.getTask、ArrayBlockingQueue.take 或 LockSupport.park,通常表示 worker 空闲等任务;停在 JDBC、HTTP 客户端、业务锁或序列化代码,才说明执行能力被具体任务占用。
beforeExecute、afterExecute 和 terminated 可用于埋点、清理上下文和生命周期观测,但钩子异常会影响 worker。指标代码必须轻量、不能反向阻塞线程池。execute 提交的 Runnable 抛未捕获异常时,worker 可能结束并由池补充;submit 把异常封装进 Future,若从不 get 或不在 afterExecute 解包,失败可能悄悄消失。
把 execute 的三个分支与回滚条件走完
execute(command) 的决定顺序定义了压力形态。第一步在 worker 数小于 core 时尝试 addWorker(command, true);失败或核心数已满后,第二步才尝试入队;入队成功不能立即返回,池状态可能刚从 RUNNING 切到 SHUTDOWN,所以还要复查,必要时把任务移除并拒绝。若入队后没有 worker,还要补一个空首任务 worker 来消费队列。只有队列不能接收时,第三步才按 maximumPoolSize 建非核心 worker;仍失败才执行拒绝策略。
复查解释了一个关闭竞态:提交线程看到 RUNNING 并成功 offer,关闭线程紧接着改变状态。没有入队后的 recheck,一个已不接受新任务的池会留下无人负责的任务。容量参数也不是独立旋钮:无界队列几乎让 maximumPoolSize 失去扩容作用;SynchronousQueue 没有存储槽,每次提交都要直接交给 worker;有界队列则显式给出排队上限与拒绝点。
worker 进入 runWorker 后反复从首任务或 getTask() 取任务。任务抛出未捕获异常时,当前 worker 异常退出,完成路径再根据池状态和最小 worker 数决定是否补建。execute 提交的异常可能到达未捕获异常处理器;submit 通常把任务包装成 FutureTask,异常被保存进 Future,调用方若从不 get,日志可能什么也没有。
cd examples/backend-development/concurrency/threadpool-queue-rejection
javac --release 17 -Xlint:all -Werror ThreadPoolFailureDemo.java
java ThreadPoolFailureDemo
# submit-returned
# future-cause=IllegalStateException第一行不是成功证据,只说明提交方法返回;第二行才从结果通道取回失败。生产封装应统一决定哪些任务必须消费 Future,哪些 fire-and-forget 任务要在包装层记录异常,哪些失败进入重试或死信。否则 completedTaskCount 会增长,业务完成量却没有对应增长。
队列选择先决定压力形态
| 队列 | 压力形态 | 常见风险 |
|---|---|---|
ArrayBlockingQueue | 明确有界,满后扩 worker 或拒绝 | 容量估错会过早拒绝或延迟过长 |
LinkedBlockingQueue(capacity) | 有界链表节点 | 额外节点内存,必须显式容量 |
无界 LinkedBlockingQueue | 几乎一直入队 | maximumPoolSize 难生效,延迟和内存累积 |
SynchronousQueue | 不存任务,提交与接收直接配对 | 慢依赖下线程数迅速扩张 |
PriorityBlockingQueue | 按优先级取任务 | 默认无界,低优先级饥饿 |
队列容量不是“能存多少”这么简单,它代表允许多少尚未执行的业务承诺。一个接口总预算 1 秒,任务平均执行 100ms,如果排队已经花了 900ms,再开始执行的价值很低。队列里的任务还可能持有请求参数、Future、上下文和大对象,容量必须同时受延迟与堆内存约束。
记录任务入队时间,在 beforeExecute 计算等待时间,比只看 queue.size() 更接近用户体验。队列长度没有增长但等待时间上升,可能是任务变慢或线程收缩;等待时间已经超过入口预算时,即使队列没满也应拒绝过期任务。
线程数要被最窄的下游约束
CPU 密集任务的并行度通常接近可用处理器数;阻塞 IO 任务可以更多,但上限不能只从 CPU 推导。数据库连接池 30、下游只允许 50 QPS,却把调用线程池开到 300,只会让更多线程堵在连接池或触发下游限流。
常用估算式:
线程数 ≈ 可用 CPU × 目标利用率 × (1 + 等待时间 / 计算时间)它只是压测起点。还要同时约束:连接池容量、下游并发/QPS、单任务内存、文件句柄、容器 CPU quota、超时和重试放大。容器内可用处理器可由 JVM 感知,必要时核对:
jcmd "$PID" VM.info | grep -i -E 'cpu|processor'
jcmd "$PID" VM.flags | grep ActiveProcessorCount
cat /sys/fs/cgroup/cpu.max 2>/dev/null || true只有在运行环境识别错误或需要稳定基线时才考虑 -XX:ActiveProcessorCount;不要用它掩盖容器资源配置问题。
四种内置拒绝策略表达四种业务后果
AbortPolicy 抛 RejectedExecutionException,适合调用方必须明确处理的核心链路。CallerRunsPolicy 由提交线程执行,能降低提交速度,但也会把延迟和工作转移到入口线程。DiscardPolicy 静默丢新任务,只适合允许丢失且有独立指标的低价值工作。
DiscardOldestPolicy 丢队列头再重试,可能扔掉最接近执行的旧任务,业务含义常常不可靠。
“支付任务不能丢”并不意味着用无界队列永远接收。进程宕机时内存队列照样丢,过期执行还可能制造重复扣款。不能丢的任务需要持久化、幂等键、Outbox 或消息系统;线程池拒绝只是应用内过载信号。
自定义策略要做的事情很少:增加拒绝计数,附上池名和饱和快照,再返回明确失败。不要在拒绝回调里做同步网络上报或无限重试,它运行在提交者线程,会把故障扩散回入口。
CallerRunsPolicy 不是免费的背压
调用者执行确实会自然减慢提交,但要看调用者是谁。HTTP 工作线程自己执行慢任务,会减少入口吞吐并抬高延迟;消息消费线程自己执行,可能延迟 offset 提交;持锁线程提交任务时被迫执行,临界区会意外拉长。
更危险的是线程池已经 SHUTDOWN 时,CallerRunsPolicy 会直接丢弃任务而不执行。使用它必须验证运行中饱和和关闭阶段两种路径,并保证有拒绝指标。
真正背压需要端到端传播:入口停止读取、返回明确过载、减少消费者拉取、缩短排队或降低并发。单个线程池策略只能影响当前提交点。
同池依赖能在还有 CPU 时把系统锁死
假设池中所有 worker 都执行父任务,父任务又向同一个池提交子任务并同步 get:
Future<Result> child = pool.submit(this::loadDetail);
return child.get();如果 worker 已被父任务占满,子任务只能排队;父任务等待子任务,worker 永不释放,形成线程饥饿死锁。CPU 可能很低,线程栈却全停在 Future.get。
解决方式可能是取消同步等待、让父任务直接执行小子任务、拆分独立池或使用结构化任务模型。仅仅扩大线程数会把死锁阈值推后,不消除依赖环。
隔离不是每个方法一个池
所有业务共用一个 asyncExecutor,导出堆积能拖慢核心查询;反过来每个类创建一个池,会制造大量空闲线程、配置碎片和无法治理的生命周期。合理隔离按故障域和资源依赖划分:
核心在线请求与低优先级后台任务分开。不同慢下游若故障不应互相传播,分开并设置各自并发上限。CPU 计算与阻塞 IO 分开,避免互相改变调度假设。
数据库任务受连接池约束,文件导出受内存/磁盘约束。
每个池必须有稳定名称、用途、负责人、容量依据、队列/拒绝语义、任务超时、关闭顺序和指标。线程池数量本身也要纳入治理,防止总线程数超过容器预算。
最小监控面要覆盖等待、执行和拒绝
System.out.printf(
"pool=%d active=%d largest=%d queue=%d task=%d completed=%d shutdown=%s%n",
pool.getPoolSize(),
pool.getActiveCount(),
pool.getLargestPoolSize(),
pool.getQueue().size(),
pool.getTaskCount(),
pool.getCompletedTaskCount(),
pool.isShutdown()
);生产指标至少包括:活跃/当前/最大线程,队列长度与容量,任务等待和执行时间,提交/完成/失败/取消/拒绝数,满载持续时间。再与入口 QPS、P99、下游耗时、连接池等待和错误率放在同一时间线。
只设置“队列达到 80% 告警”不够。一个很大的队列到 20% 时任务就可能过期;一个很小的队列瞬间满又恢复可能是预期削峰。告警应围绕等待预算、拒绝速率和持续饱和。
发布关闭要证明任务去向
先停止入口或消费者继续提交,再调用 shutdown();等待一段业务允许的 drain 时间,超时后 shutdownNow() 请求中断。返回的未开始任务必须记录、持久化补偿或明确丢弃。
任务自身还要响应中断。线程池完成了生命周期动作,不代表卡在第三方调用中的 worker 会立即退出。发布门禁应在测试中模拟:队列非空、任务运行中、任务忽略中断、线程工厂创建失败和拒绝回调异常。
线程池事故的证据顺序
jcmd "$PID" Thread.print -l > pool.txt
grep -n -E 'order-query|ThreadPoolExecutor|getTask|FutureTask|get' pool.txt看池是否 RUNNING、SHUTDOWN 或 STOP,发布动作有没有提前关闭。对照 active、queue、wait、execute、reject,而不是只看最大线程。按线程名抽样 worker 栈,找实际占用它们的依赖。
检查父任务是否等待同池子任务、锁内提交或 caller-runs。对齐下游连接池和耗时,判断是执行慢还是提交激增。先限流、降级、暂停低优先级生产者,再决定扩容;盲目加线程可能压垮下游。
上线评审模板
每个线程池至少写清:任务类型与所有者;核心/最大线程及压测依据;队列类型、容量和最大等待;拒绝后的用户/任务结果;下游容量;上下文传播;异常与取消;监控告警;关闭与补偿。代码扫描还应禁止业务裸用无界工厂池和无名称线程。
两份完整源码位于 examples/backend-development/concurrency/threadpool-queue-rejection/。ThreadPoolSaturationDemo.java 用 2 个 worker 与 1 个队列槽稳定制造拒绝;ThreadPoolFailureDemo.java 证明 submit 返回不代表任务成功。进入目录执行:
mkdir -p out
javac --release 17 -Xlint:all -Werror -d out ThreadPoolSaturationDemo.java ThreadPoolFailureDemo.java
java -cp out ThreadPoolSaturationDemo
java -cp out ThreadPoolFailureDemo稳定输出为:
pool=2 queue=1 rejected=true
submit-returned
future-cause=IllegalStateException线程池登记表不能只保存 core/max。运行门禁至少对齐提交、开始、完成、失败、取消与拒绝六个计数,并从提交/开始时间计算队列等待,从开始/结束时间计算执行耗时;关闭后还要满足“队列归零、worker 归零、未开始任务全部进入补偿或明确丢弃记录”。若 completedTaskCount 增长而业务成功量不增,优先检查 Future 结果无人消费和任务内部假成功,而不是继续扩容。
线程池的正确目标不是“永不拒绝”,而是在负载超过处理能力时,仍能用有界等待、明确失败和故障隔离保护关键路径。
可继续核对 ThreadPoolExecutor、Executor 与 BlockingQueue 的执行和内存一致性契约。
