为何df.rdd.getNumPartitions()与spark_partition_id()返回分区数不一致?
问题原因解释
这个矛盾现象的核心是**「分区数统计维度不同」**,结合代码逻辑和Spark的shuffle机制,具体原因如下:
Join后的数据分布特性
你的df_skew表的value列全为0(创建时用了[0] * 100000),而df_small的value列是0到99,两者做inner join后,只有value=0的行能匹配,最终结果的10万条数据的join key全为0。Shuffle分区的分配规则
你关闭了广播join(spark.sql.autoBroadcastJoinThreshold="-1"),Spark会采用Sort Merge Join,该Join会对两张表的join key执行shuffle操作。Spark默认的shuffle分区数spark.sql.shuffle.partitions是200,所以shuffle后会生成200个分区。但由于所有数据的join key都是0,哈希计算后会被分配到同一个分区(即结果中的191号分区),剩下的199个分区都是空的。两个分区数统计方法的差异
getNumPartitions():直接返回RDD的总分区数,不管分区是否有数据,所以会返回配置的200个分区。spark_partition_id()+groupBy:只会统计有数据的分区,空分区没有行数据,自然不会出现在统计结果中,所以只显示有10万条数据的191号分区。
内容的提问来源于stack exchange,提问作者Jerry
相关产品推荐
相关产品推荐

