Spark中parallelize分区机制及本地运行效率异常问题问询
为什么Spark本地运行时分区越多反而越慢?
这是个非常典型的Spark新手误区——直觉上觉得分区越多并行度越高就一定越快,但在本地环境+小数据集的场景下,结果往往相反。咱们结合你的代码和测试数据来拆解原因:
首先还原你的测试场景:
val longList = 1 to 1000000 toList val numsToAdd = sc.parallelize(longList,30) val st = System.nanoTime() println(numsToAdd.reduce((x,y) => x+y)) val et = System.nanoTime() println("Time spend:"+(et-st))
测试结果:
| Partitions | Time spend (nanoseconds) |
|---|---|
| default | 97916961 |
| 10 | 111907094 |
| 20 | 141691820 |
| 30 | 158264230 |
核心原因分析
1. 本地CPU核心的上限限制了并行度
Spark的并行任务数最终受限于你本地机器的CPU核心数,比如你本地是4核,那同一时间最多只能跑4个任务。就算你设置30个分区,剩下的26个任务都得排队等待,反而会增加任务调度、线程上下文切换的额外开销。而Spark的默认分区数通常会根据本地CPU核心数自动设置(一般是核心数的2-3倍),这个数量刚好能让CPU满负荷运转,同时调度开销最小。
2. 小数据集下,任务启动开销远大于计算收益
你的数据集是100万条整数,单条数据极小,就算分成默认分区,每个分区的数据量也足够小,计算耗时几乎可以忽略。而Spark每个分区对应一个任务,每个任务都需要经历启动、数据序列化/反序列化、任务上下文初始化这些固定开销。分区越多,这些固定开销的总和就越大,最终超过了并行计算带来的收益——毕竟处理3万条整数的时间,可能还不如启动一个任务的时间长。
3. Reduce操作的汇总开销随分区数增加而上升
虽然reduce会先在每个分区做局部聚合,但分区越多,最终需要汇总的中间结果就越多。比如30个分区会产生30个局部和,而默认分区可能只有4-8个局部和,汇总这些结果时需要的合并操作更多,也会带来额外的微小开销。
给你的建议
- 本地调试时,分区数不用刻意调大:设置成和CPU核心数相当或略高(比如核心数的2倍)就足够,避免不必要的调度开销;
- 大数据集才是分区发挥作用的场景:当你的数据集大到单个分区需要几十秒甚至几分钟处理时,增加分区让多个核心同时工作,才能真正提升效率。你可以试试把数据集放大到10亿条,或者每个元素换成更大的对象,再测试分区数的影响;
- 分区数的合理范围:一般来说,生产环境中分区数设置为CPU核心数的2-3倍,或者让每个分区的数据量保持在100MB-1GB之间(根据数据类型调整),是比较合理的选择。
内容的提问来源于stack exchange,提问作者Mandroid
相关产品推荐
相关产品推荐

