ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

一行parallelStream搞挂服务器:ForkJoinPool公共池的坑与自救指南

一行parallelStream搞挂服务器:ForkJoinPool公共池的坑与自救指南 先说个我自己踩过的坑。去年做一个订单批处理重构有一块逻辑要对十几万条记录做去重和汇总原来用 for 循环跑大概要 8 秒。提效嘛我看了一眼就把普通 Stream 换成了 parallelStream本地一跑2 秒完成完美。代码评审的时候大家也没多想并行两个字听着就比串行高级。结果呢上线当天晚上监控报警接口超时率从 0.1% 飙到 23%线程数曲线像心电图CPU 却一点也不高。回滚之后世界安静了。后来我蹲了半宿总算把这事的根子挖明白了。罪魁祸首就是 Java 8 并行流背后那个ForkJoinPool.commonPool()——它是整个 JVM 共享的一个线程池默认线程数等于 CPU 核数减一。你以为你在加速实际上是把所有并行任务倒进了同一口锅里锅里的水就那么多一个耗子屎就能祸害一锅汤。这篇文章我就把这事从头到尾讲透为什么一行.parallel()能把服务器搞挂公共池内部是怎么运作的出问题时怎么从线程 dump 里定位以及我后来总结出来的几条保命写法。适合正在用或准备用 Java 8 并行流的后端同学也适合面试要聊 ForkJoinPool 的看完至少能挡掉一半八股文。1. 事故还原一行 .parallel() 引发的连锁故障1.1 从加了个并行到接口超时报警那次的业务场景很简单每天晚上有一个定时任务把当天产生的订单记录捞出来按店铺、按商品、按用户维度做汇总结果写到报表表里。订单量从之前的一天几万条涨到了几十万条串行处理确实有点慢。改动也很标准// 改造前 ListOrder orders orderMapper.selectByDate(today); orders.stream() .filter(Order::isValid) .forEach(this::aggregate); // 改造后 ListOrder orders orderMapper.selectByDate(today); orders.parallelStream() .filter(Order::isValid) .forEach(this::aggregate);本地测试数据量小看不出差别到了预发环境几十万条数据确实从 8 秒降到了 2 秒一切看起来都很美好。上线之后最开始是几个下游报表数据延迟随后是订单查询接口 P99 从 200ms 涨到 3 秒多最后连完全不沾这个定时任务的其他业务也一起超时。第一反应是数据库出了问题但看慢日志和连接数一切正常看 GC 日志也没有 Full GC 的痕迹。最诡异的是 CPU。通常高并发导致超时的场景CPU 应该被打满但那次 CPU 利用率只有 20% 左右反而是线程数一路飙升活跃线程从 200 多涨到 1900 多。这种 CPU 不高、线程暴涨的形态第一反应就是有大量线程阻塞在 IO 上。1.2 jstack 里的线索公共池线程集体阻塞定位过程也不复杂一台一台机器jstack抓线程快照然后 grepcommonPool。结果让我印象深刻jstack 12345 dump.log grep -A 20 ForkJoinPool.commonPool-worker dump.log抓出来的线程栈里ForkJoinPool.commonPool-worker-1到worker-7七个线程几乎全部是WAITING状态堵在 JDBC 连接获取、外部接口调用的 socket 读上其中还有一段业务代码非常眼熟——就是我们定时任务里那个aggregate方法。ForkJoinPool.commonPool-worker-3 #16 daemon prio5 ... java.lang.Thread.State: WAITING (parking) at sun.misc.Unsafe.park(Native Method) at java.util.concurrent.locks.LockSupport.park(...) at java.util.concurrent.CompletableFuture.get(...) at com.example.report.OrderAggregator.aggregate(OrderAggregator.java:88) ...更麻烦的是除了我们的定时任务还有另外两个业务模块也用了parallelStream它们堆栈里的任务全排在后面等待。换句话说我们在公共池里提交了带阻塞 IO 的任务把 7 个线程全拖住了其他依赖公共池的并行任务全部排队排队越久接口耗时越长最后连锁引发大面积超时。这个事故链条其实特别简单公共池是全局共享的你用 7 个线程做阻塞任务等于占了 7 个公共车位后面的车全都进不来。提示排查这类问题先jstack找ForkJoinPool.commonPool-worker-*线程看它们的Thread.State和堆栈。如果大量 worker 是WAITING或TIMED_WAITING且堆栈里是 JDBC、HTTP、锁、sleep 之类的操作基本可以断定公共池被阻塞任务污染了。2. ForkJoinPool 公共池的底细为什么默认线程数是个陷阱2.1 并行流和 ForkJoinPool 的隐式绑定Java 8 引入并行流的时候官方文档写得轻描淡写parallelStream()会利用 ForkJoin 框架提高处理速度。很多人就把它当成自动多线程完全不知道底层用的是哪个线程池。答案是ForkJoinPool.commonPool()。这是一个 JVM 级的静态单例进程里所有parallelStream默认都提交给它另外CompletableFuture.supplyAsync(Supplier)不传线程池时默认用的也是它。甚至一些第三方库做异步任务如果不显式指定 Executor也会回退到它身上。更关键的是它是惰性初始化的第一次用到公共池时才创建线程。默认并行度用下面这个公式算int parallelism Math.max(1, Runtime.getRuntime().availableProcessors() - 1);为什么减一因为设计者默认调用方线程也会参与任务执行。你要处理 10 万个元素公共池 7 个 worker 在跑你当前这个业务线程也会贡献一份力凑成一个接近 CPU 核数的并行度。这个设计单看没问题但它忽略了一个事实公共池是全 JVM 共享的不是某个业务独享的而且不是所有任务都适合参与这种减一的并发度假设。公共池里的线程都是 daemon 线程名字统一是ForkJoinPool.commonPool-worker-N这也是排查时一眼认出的关键标识。2.2 工作窃取算法的优势与代价公共池的核心调度方式是工作窃取Work Stealing。传统的线程池是共享一个任务队列所有 worker 从队列里取任务队列是竞争热点。ForkJoinPool 改成了每个 worker 一个双端队列每个线程优先处理自己队列里的任务处理完了就去别的线程队列偷任务来做。这个设计对可递归拆分、CPU 密集型的任务效果非常好因为并行流的底层正是利用 Spliterator 把数据集不断二分拆成一颗任务树再分散到各个 worker 队列里。谁空闲谁就从别人队尾偷几个子任务来做这样能最大程度压榨每个线程的计算能力。但好处也会变成负担。我自己总结了三句话任务拆分本身有开销数据太少不划算。队列偷取和任务调度靠的是线程间协作如果任务里有阻塞操作等 IO 的线程不会主动让出位置偷取机制反而帮不上忙。既然是全局共享池所有使用它的业务都在互相竞争你无法控制别人往里塞什么任务。很多人以为 ForkJoinPool 线程数少所以不会把服务器搞挂。恰恰相反7 个线程不算多但都能同时阻塞住而且这个全局池一旦被阻塞任务占满影响面不是某一个接口而是整个 JVM 里所有用公共池的功能。它就像公司里唯一的公共打印机打印大文件的占着不放其他人全排队。2.3 默认并行度的两种环境陷阱第一个陷阱是容器的 CPU 感知问题。Java 8 在 Docker 容器里长期有个老大难Runtime.getRuntime().availableProcessors()在早期 JDK 上看到的是宿主机核数不是容器限制的核数。宿主机 32 核容器限制 4 核公共池默认能开出 31 个 worker 线程。线程数量超过真实可用的 CPU 之后大量时间花在线程切换上性能反而更差。JDK 8u191 之后引入了-XX:ActiveProcessorCount和默认启用的容器支持如果你还在用比较老的 JDK 8 小版本跑容器环境建议升级或者显式指定java -XX:ActiveProcessorCount4 -jar app.jar第二个陷阱是单核或并行度被调到 1 的场景。并行度变成 1 意味着公共池里可能只有一个 worker此时如果代码里出现嵌套并行流任务就非常容易互等后面专门讲。公共池还提供了一些全局参数通过-Djava.util.concurrent.ForkJoinPool.common.parallelismN可以调并行度。但我要强调的是这是全局参数调它等于影响整个 JVM 里所有使用公共池的框架和业务不到万不得已不要动。3. 服务器被搞挂的三条链路阻塞、嵌套与 IO 混用3.1 阻塞任务占满公共池7 个人都在等快递先看最典型的场景某业务为了并发查得快在parallelStream里做数据库查询、远程接口调用、Redis 读取、sleep这类任务本质是 IO 密集型。IO 密集任务的特点是 CPU 几乎不工作线程都在等网络往返。假设公共池 7 个线程一个请求过来把 5000 个任务丢给公共池。7 个 worker 各自拉起一个任务然后这个任务里需要查数据库、等 100ms 返回。在等待的 100ms 里这 7 个 worker 全在阻塞。外部接口如果慢变成 500ms、1s、超时 3s公共池就被长时间占满。后面再来的任何并行流任务都只能在队列里排队。这个排队时间取决于前面阻塞任务的总耗时跟你的业务优先级、接口紧急程度毫无关系。更麻烦的是这个池不是只有你的业务在用其他团队的代码、Spring 之类的框架异步逻辑、CompletableFuture的默认路由全被连坐。我用一个生活化的比喻你找了一个 7 人小团队干活这 7 个人全部去快递站等包裹了等多久不一定。此时客户问报表什么时候能出答案是等快递到了再说。这个比喻糙但底层逻辑一模一样。3.2 嵌套并行流外层等内层内层等线程第二条链路稍微隐蔽一点。假设你有这样一段代码外层并行流里又出现了内层并行流ListShop shops shopService.getAll(); shops.parallelStream() .forEach(shop - { ListSku skus skuService.getByShop(shop.getId()); skus.parallelStream() .forEach(this::processSku); });公共池里一共就 7 个 worker外层任务把 7 个 worker 全占住了。每个 worker 执行到内层parallelStream时需要等待内层子任务完成。内层任务去哪了排到了公共池的任务队列里。问题来了worker 们都在等内层任务的返回没人有空去执行内层任务。ForkJoinPool 有 work-stealing 机制一个 worker 在join等待时其实会去执行队列里的其他任务所以不一定立刻死锁。但一旦并行度很小比如容器里被识别成单核并行度直接变 1公共池只有一个 worker这个唯一的 worker 在外层任务里join等待内层任务完成而内层任务排到队列里需要另一个线程来执行唯一 worker 又在等待这就形成了真正的死等。就算并行度正常嵌套并行流也极度危险。任务数是指数级的外层 1000 个店铺每个店铺又有 500 个 SKU内层被拆成几百个小任务整个池的任务队列瞬间暴涨。线程栈里会看到一堆ForkJoinTask.join的等待线程不释放、任务越堆越多GC 压力和内存水位也跟着上来。结论业务代码里不要写嵌套并行流哪怕是间接嵌套也不行。什么叫间接嵌套就是你在parallelStream里调用了一个工具类方法那个方法内部自己用了parallelStream。这类藏得很深排查起来特别费劲。3.3 IO 任务与计算任务抢同一口锅第三条链路是任务类型混跑。ForkJoinPool 官方文档其实写得很清楚它适合计算密集型任务不适合在任务执行期间可能阻塞的任务。阻塞形态包括锁、IO、sleep、等待其他线程的结果。为什么不适合阻塞任务因为如果一个任务里 90% 时间在等外部结果这个 worker 线程的利用率极低但线程资源已经被占住了。如果有 1000 个任务其中 990 个是计算任务、10 个是阻塞任务理想情况应该是计算任务快速跑完阻塞任务慢慢等。但公共池只有一个10 个阻塞任务可以让大部分 worker 停摆计算任务也被迫排队。我见过一个更夸张的案例某个接口把一批数据parallelStream后每个元素里都要调用另一个服务那个服务超时时间是 30 秒。结果公共池 31 个线程全部扎进 30 秒的等待里期间这个 JVM 上其他业务的并行流全部崩溃。后来运维在配置中心里加了个开关紧急把这段逻辑改成串行分支服务才缓过来。除了阻塞任务高并发大流量本身也是问题。一个正常请求触发parallelStream处理 10 万元素时公共池要往队列里注入大量 ForkJoinTask数量级远不止 10 万次任务拆分每个请求都这么干几百个并发请求同时涌入任务队列里瞬间堆几十万、上百万个任务对象。虽然公共池有队列但它是无界的扛不住这种瞬时流量内存和 GC 直接被打爆。服务器 CPU 不高线程也不高但老年代一直在涨这也是一种常见的搞挂形态。4. 用数据判断你的任务到底适不适合并行化4.1 并行化不是免费的拆分、排队、合并都是成本别把并行当白嫖。一个串行任务丢给并行流背后至少发生这些事成本项说明拆分成本Spliterator.trySplit()将数据集不断二分每次拆分都有对象创建和数组拷贝调度成本子任务提交到 worker 队列需要处理锁、信号量、队列指针移动窃取成本空闲线程从其他队列 steal 任务涉及跨线程访问和同步合并成本子任务结果要逐层join汇总上下文切换worker 线程比较多时CPU 切换开销不可忽略单条数据处理只要几十纳秒的话上面任何一项成本都比任务本身高。数据量小的情况下并行纯粹是倒贴钱。4.2 一张评估清单 一个自测方法我每次看代码里出现parallelStream都会在心里过一遍这个清单任务类型是纯 CPU 计算还是包含 DB/HTTP/锁/sleep 等阻塞操作数据规模处理的数据量有没有到万级以上单条耗时单条数据处理时间有没有超过 100 微秒是否独立任务之间有共享状态吗有顺序依赖吗并发场景当前请求/任务的并发量有多大公共池会不会被多个业务同时争抢池的隔离这一个池是全 JVM 共享的你确认别人不往里塞阻塞任务吗自测方法也很简单拿真实数据量在真实业务环境里跑对比long start System.nanoTime(); list.stream().forEach(this::cpuTask); long serialCost System.nanoTime() - start; start System.nanoTime(); list.parallelStream().forEach(this::cpuTask); long parallelCost System.nanoTime() - start; System.out.printf(serial%d ms, parallel%d ms%n, serialCost / 1_000_000, parallelCost / 1_000_000);注意不要在单元测试里空跑一次就下结论。JIT 预热、业务数据分布、并发请求相互影响都要考虑。跑至少 5 轮看中位数和 P99而不是平均值。只有当并行后的耗时显著优于串行比如提升超过 30%才有必要引入并行流。4.3 我的经验和阈值参考我自己的经验值是这样的元素数量小于 1000肯定串行。单条处理时间小于几微秒并行收益极小拆分开销可能大于任务本身。数据量上万、单条处理时间几百微秒到几毫秒、且是纯 CPU 计算可以考虑并行但要测试。涉及网络 IO、数据库操作、分布式调用直接放弃parallelStream用显式线程池 CompletableFuture并控制并发数。一句话概括并行流解决的问题是很多 CPU 密集的独立小任务的加速问题它解决不了很多 IO 等待的问题后者需要的是更多并发承载线程和下游容量评估。5. 安全改造方案隔离、限流与替代策略5.1 Java 8 下自定义 ForkJoinPool 的局限性很多人踩过这个坑想给并行流单独建一个池于是搜到网上这么写ForkJoinPool customPool new ForkJoinPool(20); customPool.submit(() - list.parallelStream().forEach(this::process) ).join();看起来很合理定义 20 个线程的自定义 ForkJoinPool把并行流任务丢进去。但在Java 8 里这个写法多半没用并行流的内部调度逻辑仍然走ForkJoinPool.commonPool()任务根本不会跑到你自定义的池里只是把外层包装任务丢给了自定义池。这个问题在 JDK-8076446 里挂了很久Java 10 才修正。如果你生产环境就是 Java 8记住包一层自定义 ForkJoinPool 救不了 Java 8 的 parallelStream。真正要在 Java 8 里用自定义 ForkJoinPool需要绕开parallelStream自己写RecursiveTask或RecursiveAction。举个简化例子ForkJoinPool pool new ForkJoinPool(24); public class DealBatchTask extends RecursiveTaskInteger { private static final int THRESHOLD 500; private final ListOrder orders; public DealBatchTask(ListOrder orders) { this.orders orders; } Override protected Integer compute() { int size orders.size(); if (size THRESHOLD) { return orders.stream() .mapToInt(order - deal(order)) .sum(); } int mid size / 2; DealBatchTask left new DealBatchTask(orders.subList(0, mid)); DealBatchTask right new DealBatchTask(orders.subList(mid, size)); left.fork(); int rightResult right.compute(); return left.join() rightResult; } }然后通过pool.invoke(new DealBatchTask(orderList))执行。这样是能控制实际线程池但代码复杂度明显上升还得自己设计任务拆分阈值。对于绝大多数业务场景我更推荐用CompletableFuture 显式线程池既简单又可控。5.2 推荐的替代方案CompletableFuture 显式线程池举一个实际改造的例子。假设你已经确定任务里有外部调用要并发处理一批用户目标并发度 16ExecutorService bizPool new ThreadPoolExecutor( 16, 16, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue(2000), r - { Thread t new Thread(r, biz-user-handler- ThreadLocalRandom.current().nextInt(1000)); t.setDaemon(true); return t; }, new ThreadPoolExecutor.CallerRunsPolicy() ); ListCompletableFutureResult futures users.stream() .map(user - CompletableFuture.supplyAsync(() - process(user), bizPool)) .collect(Collectors.toList()); for (CompletableFutureResult future : futures) { Result result future.get(3, TimeUnit.SECONDS); }这里有几个我强烈建议养成的习惯线程名带业务含义jstack 里一眼能看出是哪个业务。队列必须是有界队列并设置拒绝策略。CallerRunsPolicy能让提交任务的线程自己执行起到天然限流效果。future.get一定要带超时时间防止个别任务卡死导致整个链路挂住。线程数不要拍脑袋定高设置前先对下游接口的容量做压测。对比一下两种方案的差异维度parallelStream 公共池CompletableFuture 显式线程池线程来源全局共享 ForkJoinPool独立线程池线程数控制默认 CPU-1全局影响大自己定义可调可控任务类型适合 CPU 密集计算、IO 都适合超时控制较难每个 future 可设置超时异常处理比较分散可统一 catch也方便 exceptionally排查难度都是 commonPool-worker难以定位业务线程名自定义一目了然5.3 保留并行流的兜底手段Semaphore 限流与分析参数如果某个场景经过测试确实适合并行流但你又担心它在公共池里影响别人可以用Semaphore给入口限个流Semaphore gate new Semaphore(4); list.parallelStream().forEach(item - { gate.acquireUninterruptibly(); try { process(item); } finally { gate.release(); } });这个做法的意义是同时最多只有 4 个子任务进入业务处理避免海量任务瞬间把公共池的任务队列塞爆。但它不能消除阻塞任务对公共池线程的占用而且并行流内部的多个子任务可能同时走到acquire上实际效果需要压测确认。它更像是一个止血带不是我日常推荐的做法。还有一些全局参数可以在特殊情况下应急比如-Djava.util.concurrent.ForkJoinPool.common.parallelism4 -Djava.util.concurrent.ForkJoinPool.common.threadFactory... -Djava.util.concurrent.ForkJoinPool.common.exceptionHandler...但再次强调这是全局配置动它之前想清楚后果。我见过有人为了救一个接口的并行流问题直接把全局公共池调大结果其他业务的任务调度全部乱掉。合理顺序是先改业务代码做线程池隔离再谈全局参数。6. 从事故沉淀出的并行流使用纪律6.1 适用边界什么情况我才敢用并行流吃过一次大亏之后我现在敢在自己负责的代码里使用并行流的场景可以用一个很窄的标准来概括数据量确实大、单条处理是纯本地 CPU 计算、任务完全独立、没有共享可变状态、调用频率可控并且经过压测证明了收益。比如一个图片灰度化处理、一个计算密集型特征提取、一次内存中大批量数据的状态判断这些任务如果单条耗时在微秒到毫秒级parallelStream是个合适的工具。但只要任务里出现数据库查询、外部 HTTP 调用、分布式缓存访问、Thread.sleep、锁等待、CompletableFuture等待嵌套我一律放弃并行流。这不是怂是算过账。6.2 Code Review 检查点这次事故之后我在团队 code review 里加了一些固定的检查项写出来给大家参考看到.parallelStream()或.parallel()必须要求作者给出性能测试数据而不是感觉更快。看到CompletableFuture.supplyAsync/runAsync没有传 Executor必须补齐线程池。看到自定义线程池检查线程名、队列长度、拒绝策略、核心线程数和最大线程数是否一致。看到并行流里出现 IO、锁、sleep、嵌套方法调用直接打回。任何全局线程池相关参数必须走配置变更评审不能悄悄写进启动脚本。我还贴了一个高危代码信号清单基本每次开会都会提一遍stream().parallel() parallelStream() ForkJoinPool.commonPool() CompletableFuture.supplyAsync(() - ...) // 没传 executor Executors.newFixedThreadPool() // 没命名没边界这些信号单独出现不一定有问题但出现在一起爆炸概率会指数级上升。6.3 监控与应急如何提前发现公共池被打满事后我养成了把程序化采样加进监控的习惯定期输出公共池指标ForkJoinPool pool ForkJoinPool.commonPool(); log.info( commonPool - active{}, running{}, queuedTasks{}, steals{}, pool.getActiveThreadCount(), pool.getRunningThreadCount(), pool.getQueuedTaskCount(), pool.getStealCount() );如果activeThreadCount长期等于并行度queuedTaskCount持续增长说明公共池已经在排队了。再配合jstack看commonPool-worker-N的状态就能判断到底是哪个业务的阻塞任务在里面捣乱。应急处理我有一套固定动作先看 jstack确认公共池线程堆栈是否指向自己团队的业务。如果指向立刻切到串行分支或关掉并行逻辑。如果公共池被其他业务占满走熔断降级保证核心接口不受影响。临时调大common.parallelism只能作为 5 分钟级别的止血操作不能作为长期方案。事后必须要求责任人提交根因分析把所有使用公共池的代码列出来一条条过。我自己最大的体会是并行流这层语法糖太甜了甜到让人忘了它背后是一个全局共享、线程数极少、还不适合阻塞任务的线程池。代码里写.parallel()只有一瞬间的爽但线上排查的苦往往是半夜三更加倍还回来的。如果你团队里还有人喜欢随手加parallel()把这篇文章转给他最好再拉一个 code review 会议当场把公共池打满的链路讲一遍比看多少文档都管用。
返回列表