如何将PySpark各列空值计数结果整合为结构化DataFrame
PySpark空值统计结果整合方案
PySpark原生实现(推荐)
不需要单独执行三次统计查询,可直接输出结构化的空值统计结果,适合大数据量表场景,计算全程分布式执行,无需拉取数据到本地:
-- 直接执行SQL即可得到目标DataFrame SELECT column_name, null_count FROM ( SELECT COUNT(CASE WHEN student_id IS NULL THEN 1 END) AS student_id, COUNT(CASE WHEN student_scores IS NULL THEN 1 END) AS student_scores, COUNT(CASE WHEN student_health IS NULL THEN 1 END) AS student_health FROM student_table ) t LATERAL VIEW EXPLODE(MAP( 'student_id', student_id, 'student_scores', student_scores, 'student_health', student_health )) tmp AS column_name, null_count
如果需要适配表的所有字段,不需要手动写列名,可以用Python动态生成查询逻辑:
# 获取表的全部字段列表 table_columns = [col.name for col in spark.table("student_table").schema] # 生成空值统计逻辑片段 count_expr = ",".join([f"COUNT(CASE WHEN `{c}` IS NULL THEN 1 END) AS `{c}`" for c in table_columns]) # 生成列转行逻辑片段 map_expr = ",".join([f"'{c}', `{c}`" for c in table_columns]) # 执行查询得到结果DataFrame null_count_df = spark.sql(f""" SELECT column_name, null_count FROM (SELECT {count_expr} FROM student_table) t LATERAL VIEW EXPLODE(MAP({map_expr})) tmp AS column_name, null_count """) # 查看统计结果 null_count_df.show()
Pandas实现
如果你已经完成了三次单独查询,想要快速合并结果,可以用Pandas构造目标DataFrame:
import pandas as pd # 提取三次查询的统计结果 id_null = spark.sql("select count(*) from student_table where student_id is NULL").collect()[0][0] score_null = spark.sql("select count(*) from student_table where student_scores is NULL").collect()[0][0] health_null = spark.sql("select count(*) from student_table where student_health is NULL").collect()[0][0] # 构造Pandas格式的统计结果 result_pd = pd.DataFrame({ "column_name": ["student_id", "student_scores", "student_health"], "null_count": [id_null, score_null, health_null] }) # 如有需要可以转回PySpark DataFrame result_spark_df = spark.createDataFrame(result_pd)
两种方案输出的结果均为每行对应一个字段的空值统计,包含column_name(字段名)和null_count(空值数量)两个字段,符合预期输出格式。
内容的提问来源于stack exchange,提问作者Jaehyeok Kwak
相关产品推荐
相关产品推荐

