You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

请求协助:为Databricks数据库表列DataFrame添加空值统计列

为Databricks表结构DataFrame新增空值统计列

现有表结构DataFrame

DatabaseTableColumnColumnType
defaulttable1column1string
defaulttable1column2boolean
defaulttable2column3integer
defaulttable2column4string
defaulttable2column5string

预期结果DataFrame

DatabaseTableColumnColumnTypeNullsPercentage
defaulttable1column1string3455%
defaulttable1column2boolean00%
defaulttable2column3integer9876090%
defaulttable2column4string5672152%
defaulttable2column5string15121%

完整实现代码

假设你已有的表结构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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 06:20:25