如何基于合并后的规则校验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"} ]
修改思路
- 区分规则类型:遍历规则时,通过判断键名是
column(单列)还是columns(多列)执行不同统计逻辑 - 单列规则:保留原有逻辑,统计目标列中匹配值的行数占总行数的比例
- 多列规则:统计所有指定列中匹配值的单元格总数,除以这些列的总单元格数(总行数×列数)得到占比;同时用聚合操作替代多次
count()调用,提升效率 - 边界处理:添加除数为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
相关产品推荐
相关产品推荐

