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

如何基于合并后的规则校验PySpark DataFrame多列数据质量?

适配多列规则的PySpark数据占比统计修改方案

问题说明

现有PySpark测试DataFrame,需将last_name和first_name列的"NA"值占比规则合并为多列形式,同时保留country列"USA"值的独立单列规则。原有代码仅支持单列规则,需修改以适配新规则格式,实现对应统计需求。

原单列规则示例:

rules = [
    {"column": "last_name", "value": "NA", "name": "Percentage of 'NA' Values in Last Name"},
    {"column": "first_name", "value": "NA", "name": "Percentage of 'NA' Values in First Name"}
]

混合单/多列的新规则示例:

rules = [
    {"columns": ["last_name", "first_name"], "value": "NA", "name": "Percentage of 'NA' Values in Last Name and First Name"},
    {"column": "country", "value": "USA", "name": "Percentage of 'USA' Values in Country"}
]

修改思路

  1. 区分规则类型:遍历规则时,通过判断键名是column(单列)还是columns(多列)执行不同统计逻辑
  2. 单列规则:保留原有逻辑,统计目标列中匹配值的行数占总行数的比例
  3. 多列规则:统计所有指定列中匹配值的单元格总数,除以这些列的总单元格数(总行数×列数)得到占比;同时用聚合操作替代多次count()调用,提升效率
  4. 边界处理:添加除数为0的判断,避免计算报错

修改后的完整代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as spark_sum

# 初始化SparkSession(示例)
spark = SparkSession.builder.appName("ColumnPercentageCheck").getOrCreate()

# 示例DataFrame(可替换为你的实际数据)
data = [
    ("NA", "Smith", "USA"),
    ("John", "NA", "Canada"),
    ("NA", "NA", "USA"),
    ("Jane", "Doe", "USA")
]
df = spark.createDataFrame(data, ["last_name", "first_name", "country"])

# 新规则
rules = [
    {"columns": ["last_name", "first_name"], "value": "NA", "name": "Percentage of 'NA' Values in Last Name and First Name"},
    {"column": "country", "value": "USA", "name": "Percentage of 'USA' Values in Country"}
]

percentages = []
total_count = df.count()  # 仅计算一次总行数,提升效率

for rule in rules:
    value = rule["value"]
    name = rule["name"]
    percentage = 0.0

    if "columns" in rule:
        # 处理多列规则
        columns = rule["columns"]
        # 构建表达式:统计每列中等于目标值的行数
        count_exprs = [spark_sum((col(c) == value).cast("int")).alias(f"{c}_count") for c in columns]
        # 一次性聚合所有列的统计结果
        count_results = df.agg(*count_exprs).collect()[0]
        # 计算所有列的目标值总数
        total_matches = sum(count_results[col_name] for col_name in count_results)
        # 计算总单元格数
        total_cells = total_count * len(columns)
        if total_cells > 0:
            percentage = (total_matches / total_cells) * 100
    else:
        # 处理单列规则
        column = rule["column"]
        match_count = df.filter(col(column) == value).count()
        if total_count > 0:
            percentage = (match_count / total_count) * 100

    percentages.append({"name": name, "percentage": percentage})

# 输出统计结果
for result in percentages:
    print("{}: {:.2f}%".format(result["name"], result["percentage"]))

代码说明

  • 多列统计优化:使用agg函数一次性聚合所有列的匹配行数,避免多次遍历DataFrame,大幅提升性能
  • 边界处理:当总行数或总单元格数为0时,默认占比为0,避免除以0错误
  • 兼容性:同时支持原有单列规则和新增的多列规则,无需修改原有规则结构

内容的提问来源于stack exchange,提问作者PythonLearner

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 18:42:51