多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

11-Fork/Join 与工作窃取:并行流为什么能跑满 CPU,又有哪些坑

11-Fork/Join 与工作窃取:并行流为什么能跑满 CPU,又有哪些坑 for循环处理 1000 万条数据单线程要跑 10 秒list.parallelStream()一行改成并发有时候 2 秒就完。魔法背后是Fork/Join 框架和工作窃取Work-Stealing。这一篇讲清它怎么把大任务拆小、怎么让空闲线程去偷别人的活以及并行流那些「用了反而更慢」的隐藏陷阱。一、Fork/Join 是什么Fork/Join 是 JDK 7 引入的并行框架专为「可以递归拆分成小块的 CPU 密集型任务」设计。核心思想Fork分解把大任务递归拆成足够小的子任务。Join合并等所有子任务算完把结果合并。它不是一个普通线程池而是ExecutorService的一个特殊实现ForkJoinPool用工作窃取算法把 CPU 吃满。二、工作窃取让核不闲着普通线程池的任务是「一个共享队列所有线程来抢」容易有竞争且一个线程干完就闲着。Fork/Join 反过来每个工作线程有自己的双端队列deque自己从头取任务干。当某线程做完了自己的队列线程 A 队列: [t1 t2 t3 t4 t5] ← A 从头部取干得快 线程 B 队列: [t6 t7 t8] ← B 慢 A 干完自己后从 B 队列的【尾部】偷一个任务来干「偷」发生在别人队列的另一端和主人取任务的端错开极大减少锁竞争。结果快的线程去帮慢的所有核尽量都在忙负载均衡自动达成。这正是它能「跑满 CPU」的原因。三、 RecursiveTask / RecursiveAction自己写 Fork/Join 任务继承这两个有返回值用RecursiveTask无返回值用RecursiveAction。关键是「拆到多小才停止」——阈值太小任务太多 overhead 大太大并行度不够。常见按数据量定阈值。classSumTaskextendsRecursiveTaskLong{staticfinalintTHRESHOLD1000;finallong[]arr;finalintlo,hi;SumTask(long[]a,intlo,inthi){arra;this.lolo;this.hihi;}protectedLongcompute(){if(hi-loTHRESHOLD){// 足够小直接算longs0;for(intilo;ihi;i)sarr[i];returns;}intmid(lohi)1;SumTaskleftnewSumTask(arr,lo,mid);SumTaskrightnewSumTask(arr,mid,hi);left.fork();// 异步拆出左半longrright.compute();// 当前线程直接算右半不浪费longlleft.join();// 等左半结果returnlr;}}longtotalnewForkJoinPool().invoke(newSumTask(big,0,big.length));注意left.fork()后当前线程right.compute()而不是right.fork()——避免多开一个线程空等这是标准写法能少一次线程占用。四、并行流Fork/Join 的语法糖Stream.parallel()/parallelStream()就是把流水线交给公共ForkJoinPool.commonPool()跑longsumlist.parallelStream().filter(x-x0).mapToLong(x-x).sum();底层自动用 Fork/Join 把数据切片、各核并行处理、再合并。写起来一行背后是完整的工作窃取并行。五、并行流的 5 个坑重点共享可变状态 → 数据错乱并行流里千万别改外部共享变量累加要用reduce/collect返回新值别用外部total x// 错误多线程改同一个 total结果错longtotal0;list.parallelStream().forEach(x-totalx);// 正确用 reduce每个线程算自己的局部和再合并longtotallist.parallelStream().reduce(0L,Long::sum);用公共池跑 IO 阻塞commonPool默认并行度 核数-1且是守护线程、全局共享。若在并行流里做 DB/HTTP 调用阻塞会占满公共池拖垮所有用并行流的地方包括其他业务的CompletableFuture。数据量小反而更慢拆任务、合并、线程协调都有开销。几千条以下的数据并行流的 overhead 可能超过并行收益直接用串行流更快。ThreadLocal失效公共池线程是复用的、不绑定你的调用线程链ThreadLocal/InheritableThreadLocal在并行流里不按预期传递要用TL传上下文会丢。顺序敏感操作不能用forEach在并行下不保证顺序需要顺序用forEachOrdered但会牺牲并行度。findFirst在并行下也可能比串行慢要保持顺序语义。六、什么时候该用什么时候别用场景建议大数据量 CPU 密集求和/过滤/转换✅ 并行流首选数据量小 数千❌ 直接用串行任务里有 IO/网络/锁❌ 别用并行流用自定义 IO 线程池需要强顺序输出⚠️ 用forEachOrdered收益打折要控制并行度/隔离⚠️ 自己new ForkJoinPool()提交别用公共池自定义池隔离示例避免污染公共池ForkJoinPoolpoolnewForkJoinPool(8);longrpool.submit(()-list.parallelStream().reduce(0L,Long::sum)).get();七、和线程池的关系Fork/Join 是为「计算」而生的并行框架不适合长阻塞任务。日常业务「异步处理 IO」还是用普通ThreadPoolExecutor只有「可拆分的大计算」才上 Fork/Join / 并行流。两者不是替代关系是分工IO/任务调度用线程池纯计算拆并用 Fork/Join。总结Fork/Join 用「分治 工作窃取」把 CPU 吃满每个线程有自己双端队列干完去偷别人队列尾部的任务自动负载均衡。并行流是它的一行语法糖。reduce/collect返回的合并才线程安全别在并行流里改共享变量公共commonPool是全局共享的守护线程池跑 IO 阻塞会拖垮所有并行流用户小数据量、强顺序、依赖ThreadLocal时别用并行流。IO 密集型仍用普通线程池Fork/Join 专攻可拆分的大计算。
返回列表