PySpark repartition()元素集中单分区问题咨询(Spark 2.1.1)
嘿,我来帮你拆解下这个问题~你的操作其实没有错,这个看起来“异常”的结果是Spark针对小数据本地RDD的优化导致的,具体原因和验证方法如下:
为什么会出现这种现象?
你用sc.parallelize(range(20), 8)生成的是ParallelCollectionRDD,这类RDD的数据直接存储在Driver端的本地集合里。对于这种数据量极小的本地RDD,Spark会跳过分布式Shuffle的流程,直接在Driver端完成分区重划分——它不会把每个元素通过网络分发到Executor的不同分区,而是直接对原本地集合做重新分区处理,这就导致你看不到预期的元素打乱分配效果,甚至出现所有元素集中到一个分区的情况。
这种优化的初衷是减少不必要的网络开销,但确实会干扰我们测试Shuffle逻辑的预期结果。
如何验证并触发真正的Shuffle?
你可以通过对RDD做一次简单的转换操作,让它脱离ParallelCollectionRDD的类型,再调用repartition(),就能看到正常的Shuffle效果了:
# 做一次map转换,生成非本地集合类型的RDD rdd = sc.parallelize(range(20), 8).map(lambda x: x) # 再执行repartition并查看分区 print(rdd.repartition(8).glom().collect())
这时候你会发现元素被按照哈希规则分配到不同分区,符合你最初的预期。
另外,你也可以打开Spark UI查看Stages页面:如果是ParallelCollectionRDD执行repartition(),页面里不会出现Shuffle Read/Write相关的Stage;而经过转换后的RDD执行repartition(),就能看到Shuffle对应的Stage了。
补充:HashPartitioner的实际逻辑
当真正触发分布式Shuffle时,repartition()默认使用HashPartitioner。对于Int类型的元素,它会用元素本身的值作为哈希值,对目标分区数取模来分配分区——比如元素0会分到0 % numPartitions对应的分区,元素1分到1 % numPartitions对应的分区,以此类推。
如果你需要测试repartition()的Shuffle效果,建议使用经过转换的RDD或者从外部数据源(比如文件)读取的RDD,这样就能避开Spark的本地优化,得到符合预期的结果啦~
内容的提问来源于stack exchange,提问作者xcynn

