请求协助:为Databricks数据库表列DataFrame添加空值统计列
为Databricks表结构DataFrame新增空值统计列
现有表结构DataFrame
| Database | Table | Column | ColumnType |
|---|---|---|---|
| default | table1 | column1 | string |
| default | table1 | column2 | boolean |
| default | table2 | column3 | integer |
| default | table2 | column4 | string |
| default | table2 | column5 | string |
预期结果DataFrame
| Database | Table | Column | ColumnType | Nulls | Percentage |
|---|---|---|---|---|---|
| default | table1 | column1 | string | 345 | 5% |
| default | table1 | column2 | boolean | 0 | 0% |
| default | table2 | column3 | integer | 98760 | 90% |
| default | table2 | column4 | string | 56721 | 52% |
| default | table2 | column5 | string | 1512 | 1% |
完整实现代码
假设你已有的表结构DataFrame名为 schema_df,以下是完整的Python实现:
from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, LongType # 提取所有需要处理的数据库和表组合 table_pairs = schema_df.select("Database", "Table").distinct().collect() # 初始化空的统计结果DataFrame null_stats_df = spark.createDataFrame([], StructType([ StructField("Database", StringType()), StructField("Table", StringType()), StructField("Column", StringType()), StructField("Nulls", LongType()), StructField("Percentage", StringType()) ])) # 逐个表计算空值统计 for pair in table_pairs: db = pair["Database"] table = pair["Table"] full_table = f"`{db}`.`{table}`" # 读取目标表数据 target_df = spark.sql(f"SELECT * FROM {full_table}") # 获取表的总行数 total_rows = target_df.count() # 计算每列的空值数量 null_count_cols = [F.count(F.when(F.col(c).isNull(), c)).alias(c) for c in target_df.columns] null_counts_wide = target_df.select(null_count_cols) # 将宽格式的统计结果转为长格式,方便后续关联 null_counts_long = null_counts_wide.select( F.explode(F.map_from_entries(F.array(*[F.struct(F.lit(c), F.col(c)) for c in target_df.columns]))).alias("col_stats") ).select( F.col("col_stats.key").alias("Column"), F.col("col_stats.value").alias("Nulls") ).withColumn("Database", F.lit(db)).withColumn("Table", F.lit(table)) # 计算空值占比,处理总行数为0的情况避免报错 null_stats = null_counts_long.withColumn( "Percentage", F.when(total_rows == 0, "0%") .otherwise(F.concat(F.round((F.col("Nulls") / total_rows) * 100, 0).cast("string"), "%")) ) # 合并到统计结果DataFrame null_stats_df = null_stats_df.union(null_stats) # 将统计结果与原始表结构DataFrame关联,得到最终结果 final_result_df = schema_df.join(null_stats_df, on=["Database", "Table", "Column"], how="inner") # 展示最终结果 final_result_df.show()
关键说明
- 遍历所有唯一的数据库表组合,确保每个表只处理一次
- 计算表的总行数,用于空值占比的计算
- 将宽格式的空值统计结果转为长格式,便于和原始表结构表进行关联
- 处理总行数为0的边界情况,避免出现除以0的错误
- 通过关联操作将空值统计信息合并到原始表结构DataFrame中,得到包含所有字段的最终结果
内容的提问来源于stack exchange,提问作者Guillermo Colom
相关产品推荐
相关产品推荐

