如何在Spark中使用DataFrame统计各分区记录数且不忽略空分区
在Spark DataFrame中统计所有分区(含空分区)的记录数
直接使用groupBy(spark_partition_id()).count()只能统计包含数据的分区,空分区不会出现在结果集中。要实现包含空分区的统计,需要先生成全部分区ID的数据集,再与原统计结果做左连接补全空值。
实现步骤:
生成全部分区ID的DataFrame
通过RDD的getNumPartitions()获取总分区数,生成从0到总分区数-1的序列并转为DataFrame:import org.apache.spark.sql.functions._ import spark.implicits._ val totalPartitions = df.rdd.getNumPartitions val allPartitionsDF = (0 until totalPartitions).toDF("partition_id")统计有数据分区的记录数
基于原DataFrame统计各分区记录数,为列名添加别名方便后续连接:val nonEmptyPartitionCountsDF = df .groupBy(spark_partition_id().alias("partition_id")) .count()左连接补全空分区记录数
将全部分区ID数据集与分区统计结果左连接,用coalesce把空分区的count值替换为0:val allPartitionCountsDF = allPartitionsDF .join(nonEmptyPartitionCountsDF, Seq("partition_id"), "left_outer") .select( col("partition_id"), coalesce(col("count"), lit(0)).alias("record_count") )
执行以上代码后,allPartitionCountsDF会包含所有分区ID,空分区对应的record_count为0。
内容的提问来源于stack exchange,提问作者Mohammad Mahfooz Alam
相关产品推荐
相关产品推荐

