Spark DataFrame分区数查询结果不一致问题咨询
Spark分区数查询结果不一致问题解析
执行的操作与结果
- 查看分区ID及对应记录数的代码:
输出结果:data_frame.toDF().withColumn("partitionId", spark_partition_id()).groupBy("partitionId").count().orderBy(asc("count")).show()+-----------+-----+ |partitionId|count| +-----------+-----+ | 3| 22| +-----------+-----+ - 查看分区总数的代码:
输出结果:logger.warning('Num partitions: %s', data_frame.toDF().rdd.getNumPartitions())WARNING:root:Num partitions 4
不一致的核心原因
- 空分区不会被统计:
spark_partition_id()只会返回包含数据的分区ID,那些没有任何记录的空分区,不会出现在groupBy的统计结果里。你的数据仅在第3个分区有22条记录,剩余3个分区都是空的,所以最终只显示这一个有数据的分区。 getNumPartitions()统计所有分区:这个API是直接获取RDD的总分区数量,不管分区里有没有数据,哪怕是空分区也会被算入总数,所以返回的是实际创建的4个分区的总数。
验证空分区的方法
如果要确认空分区的存在,可以用mapPartitions遍历每个分区,输出分区ID和是否有数据:
def check_partition(iterator): pid = spark_partition_id() has_data = len(list(iterator)) > 0 yield (pid, has_data) data_frame.toDF().rdd.mapPartitions(check_partition).collect()
执行后会返回所有4个分区的ID和数据存在状态,能直接看到另外3个分区都是空的。
内容的提问来源于stack exchange,提问作者Mike
相关产品推荐
相关产品推荐

