七参数 · 反直觉的执行顺序 (core → 队列 → max → 拒绝) · 拒绝策略 · 为什么禁用 Executors 便捷工厂
线程池 = 参数化的"用工制度": 核心员工(core)常驻, 忙不过来先排队(队列), 队伍排满了才招救急工(max), 还不行就拒绝(Handler)。最反直觉也最常考的是顺序 — 先入队、后扩容; 最常出事的是队列 — 无界队列让 maximumPoolSize 形同虚设, 任务堆积到 OOM。这就是阿里规约"禁用 Executors、必须手动 new"的根本原因。
Executors.newFixedThreadPool(10); // 错: 无界队列 → 堆积 OOM new ThreadPoolExecutor(10, 20, 60, SECONDS, new ArrayBlockingQueue<>(1000), ...); // 对: 显式七参数
pool.execute(task); // 异常直接打到 UncaughtExceptionHandler Future<?> f = pool.submit(task); f.get(); // 关键: 不 get() 异常被包进 Future 静默吞掉
allowCoreThreadTimeOut(true) 可连核心一起收 — 低峰省资源, 适合流量波谷极深的长尾服务。 ThreadPoolExecutor p = new ThreadPoolExecutor( core, max, 60, SECONDS, queue); // 只回收 (core, max] p.allowCoreThreadTimeOut(true); // 关键: 连核心一起收
int n = Runtime.getRuntime().availableProcessors(); threads = n + 1; // CPU 密集 threads = n * (1 + wait / service); // IO 密集, 压测定终值
pool.setCorePoolSize(32); // 运行时生效, 无需重启 pool.setMaximumPoolSize(64); // 配套: 活跃数 / 队列水位 / 拒绝计数 → 配置中心联动
new LinkedBlockingQueue<Runnable>(); // 错: 无界, max 永不生效 new ArrayBlockingQueue<>(1000); // 对: 有界+拒绝=快速失败
ctx.set(userId); try { process(); } finally { ctx.remove(); } // 关键: 线程复用, 不 remove = 脏数据+泄漏
ThreadPoolExecutor pool = new ThreadPoolExecutor( 8, 16, // core / max: 压测起点 60, TimeUnit.SECONDS, new ArrayBlockingQueue<>(500), // 有界! 水位=缓冲突发, 不是垃圾桶 new ThreadFactoryBuilder() .setNameFormat("order-export-%d").build(), // 命名: jstack 一眼定位 new ThreadPoolExecutor.CallerRunsPolicy()); // 反压 (或自定义降级)
// 暴露到 Prometheus: gauge.set(pool.getActiveCount(), pool.getQueue().size(), rejectedCount.get()) // 告警联动配置中心 (Apollo/Nacos) 热调参, 不重启: pool.setCorePoolSize(newCore); // 运行时生效 pool.setMaximumPoolSize(newMax); // 大促前: 队列水位 >80% → 自动 core 8→32, 平峰回收
Future<?> f = pool.submit(task); try { f.get(3, TimeUnit.SECONDS); } // 不 get: 任务死了都没声 catch (ExecutionException e) { log.error("task failed", e.getCause()); // 真实业务异常在 cause 里 }
Tomcat 线程只做接入, 慢业务(导出/推送)必须独立池, 防互相拖垮:
// Tomcat: 快速接入 (server.tomcat.threads.max=200) # 业务: 慢操作独立池 — 导出任务再慢也不影响下单接口 @Bean("exportPool") public ThreadPoolExecutor exportPool() { return new ThreadPoolExecutor(4, 8, 60, TimeUnit.SECONDS, new ArrayBlockingQueue<>(200), namedThreadFactory("export"), // 命名区分: jstack 一眼归属 new ThreadPoolExecutor.AbortPolicy()); } # @Async("exportPool") 或手动 submit — 明确指定, 不用默认池
并行调 N 个下游, 按完成顺序消费而不是提交顺序等待:
ExecutorCompletionService<Quote> ecs =
new ExecutorCompletionService<>(pool);
for (Vendor v : vendors) ecs.submit(() -> v.quote(req));
for (int i = 0; i < vendors.size(); i++) {
Future<Quote> f = ecs.poll(500, MILLISECONDS); // 谁先回先拿
if (f != null) merge(f.get()); // 超时不候: 慢 vendor 不拖整体
}
声明式并发编排, 但默认池是 ForkJoinPool.commonPool — 生产必须显式传池:
CompletableFuture
.supplyAsync(() -> queryUser(uid), ioPool) // 显式业务池!
.thenCombine(
CompletableFuture.supplyAsync(() -> queryOrders(uid), ioPool),
(user, orders) -> render(user, orders))
.orTimeout(800, MILLISECONDS) // JDK9+ 整链超时
.exceptionally(e -> fallbackPage()); // 降级页兜底
自定义拒绝策略: 不丢任务, 转投消息队列削峰:
new ThreadPoolExecutor.RejectedExecutionHandler() { public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { log.warn("pool saturated, queue={}", e.getQueue().size()); rejectedCounter.inc(); // 拒绝计数: 核心告警指标 mqProducer.send("task-retry", serialize(r), delay(30s)); // 30s 后回放 } }
支付回调与营销推送共用池 = 推送风暴时支付被拖死; 按重要性分组隔离:
pools.put("payment", newPool(32, 32, 1000)); // 核心交易: 大配额 pools.put("marketing", newPool(4, 8, 200)); // 可降级: 小配额+丢弃策略 # 每个池独立监控水位与拒绝数 — 一个池打满, 告警只指向那个业务域
线程数不是猜的 — 固定负载扫描, 找吞吐拐点与队列水位的关系:
// 阶梯压测: 逐步加压, 每档记录三件事 for threads in 50 100 150 200 300; do run_load --threads=$threads --duration=120s record: QPS | P99 | pool.active + queue.size # 三线同图 done # 拐点特征: QPS 不再涨而 queue 持续涨 → 该档位就是容量上限, 参数定在拐点前 20%
定时任务池的两个必知: 异常会杀掉后续调度、任务要幂等:
ScheduledExecutorService ses = Executors.newScheduledThreadPool(4, namedThreadFactory("sync")); ses.scheduleWithFixedDelay(() -> { try { syncFromUpstream(); } // 必须 try 包住! catch (Throwable t) { log.error("sync fail", t); } // 抛出=该任务永久停摆 }, 0, 30, TimeUnit.SECONDS); // fixedDelay: 上一轮跑完再计时
Executors.newFixedThreadPool(8); // 错: 无界队列, max 永不生效 new ThreadPoolExecutor(8, 16, 60, SECONDS, new ArrayBlockingQueue<>(500), handler); // 对
// 错: 不配置 → SimpleAsyncTaskExecutor, 每任务新开线程 // 对: 实现 AsyncConfigurer 返回自定义 ThreadPoolTaskExecutor
userCtx.set(u); // 错: 用完不 remove → 复用线程读到旧用户 try { process(); } finally { userCtx.remove(); } // 对
pool.submit(task); // 错: 异常包进 Future 静默吞 pool.submit(task).get(); // 对: 或 afterExecute 统一记日志
Runtime.getRuntime().addShutdownHook(new Thread(() -> { pool.shutdown(); // 1. 不收新任务 pool.awaitTermination(30, SECONDS); // 2. 等存量跑完(示意) pool.shutdownNow(); })); // 3. 超时强制收尾
child.get(); // 错: 同池父子任务互等 → 池死锁 // 对: 父子分池, 或 CompletableFuture 组合不阻塞
for (int i = 0; i < 1_000_000; i++) pool.submit(job); // 错: 百万任务排队 = 内存炸弹 // 对: 分批 + Semaphore 限流, 队列容量匹配消费速率
(core = 16, max = 16) // 错: 突发只能排队 (core = 16, max = 32, keepAlive = 60s) // 对: 救急余量
core = 512; // 错: IO 型不看下游容量盲目拉高 → 切换风暴 vmstat 1 # 对: 看 cs 列; 线程数由下游容量与等待比决定
keepAlive = 0s; // 错: 救急线程建完即毁, 反复建线程 keepAlive = 60s; // 对: 让突发期线程活过波峰
queue = 10; // 错: 稍慢就拒绝, max 永在救火边缘 // 对: 容量 = 突发时长 × 消费速率, 配水位告警
new ThreadPoolExecutor.CallerRunsPolicy(); // 错: 反压压垮 HTTP 线程 // 对: 在线服务落 MQ / 降级; CallerRuns 只用于离线
new ThreadPoolExecutor.DiscardPolicy(); // 错: 静默丢单 (r, e) -> { log.error("rejected", r); metric.inc(); } // 对
// 错: 一处 execute 一处 submit, 异常出口不一致 // 对: 团队统一一种 + afterExecute / 强制 get 兜底
// 错: OOM 也被 submit 的 Future 吞掉, 掩盖致命问题 protected void afterExecute(Runnable r, Throwable t) { ... } // 对
pool.setMaximumPoolSize(4); // 错: 小于当前 core → IllegalArgumentException // 对: 先 setCorePoolSize 再 setMaximumPoolSize, max ≥ core
sched.scheduleAtFixedRate(() -> {
try { job(); } catch (Throwable e) { log.error(e); } // 对
}, 0, 1, MINUTES); // 错: 裸抛一次 → 后续全部静默停摆pool.submit(task); // 错: shutdown 后 → RejectedExecutionException → 500 // 对: 停机先摘流量, catch 拒绝异常转友好提示
// 错: 随手 daemon=true → JVM 退出排队任务全丢 // 对: 后台任务 daemon+显式停机等待; 关键任务非 daemon
CompletableFuture.supplyAsync(job); // 错: 默认 commonPool, 容器限核或=1 全串行 CompletableFuture.supplyAsync(job, ioPool); // 对: 显式指定池