如何优化PySpark代码:替换循环实现列空值统计
Spark高效计算DataFrame列空值计数与占比(避免循环)
原代码的核心问题是循环遍历每列并执行count(),每个count()都是Spark Action操作,会触发独立作业,导致数据被重复扫描,性能损耗极大。以下是仅触发两次Action的优化方案,输出结果与原函数完全一致:
优化实现代码
from pyspark.sql import functions as F def get_null_count_and_percentage(df): # 仅执行一次count获取总行数 total_count = df.count() # 生成所有列的空值计数聚合表达式,一次扫描完成所有列统计 agg_exprs = [ F.sum(F.when(F.col(col).isNull(), 1).otherwise(0)).alias(f"{col}_null_count") for col in df.columns ] agg_df = df.agg(*agg_exprs) # 将聚合得到的宽表转换为(列名, 空值计数)的长表结构 null_counts_df = agg_df.select( F.explode( F.array([ F.struct( F.lit(col).alias("column_name"), F.col(f"{col}_null_count").alias("null_counts") ) for col in df.columns ]) ).alias("null_info") ).select("null_info.*") # 计算空值占比并保留三位小数 result_df = null_counts_df.withColumn( "null_percentage", F.round((F.col("null_counts") / total_count) * 100, 3) ) return result_df
关键优化点
- 减少Action次数:仅执行
df.count()和最终的聚合+转换两次Action,避免循环触发多次作业 - 单次数据扫描:通过
agg()函数一次性完成所有列的空值计数统计,数据只被扫描一次 - 结构转换高效:使用
explode+array+struct将宽表转长表,操作完全基于Spark批处理逻辑,性能远高于Python循环
内容的提问来源于stack exchange,提问作者Anand Khond
相关产品推荐
相关产品推荐

