如何标记两个DataFrame中属性不匹配的错误记录?
批量对比基准与待检查DataFrame并标记错误行及属性位置
Pandas 实现方案
实现逻辑
- 定义列映射关系:明确基准DataFrame(base_df)和待检查DataFrame(check_df)中需要对比的属性对应关系,比如基准的
attribute_1对应待检查的Attribute_1。 - 关联两个DataFrame:以
my_id和parent_id作为关联键,将基准DF中需要对比的属性合并到待检查DF中,确保每行能匹配到对应的基准值。 - 逐属性生成错误标记:对每一组对应属性,判断待检查值与基准值是否一致,生成布尔型的临时标记列。
- 汇总错误属性:遍历每行的临时标记列,收集所有不一致的待检查列名,形成错误属性列表。
- 筛选并整理结果:过滤出存在错误的行,保留待检查DF的原始列和错误属性列表列,移除中间临时列。
代码实现
import pandas as pd # 构造基准DataFrame base_df = pd.DataFrame({ 'my_id': ['ABC', 'ABS', 'ABB', 'ABB', 'APP'], 'parent_id': ['DEF', 'DES', 'DEG', 'DEG', 'DRE'], 'attribute_1': ['A-', 'A-', 'A', 'B-', 'C-'], 'attribute_2': [378.8, 388.8, 908.8, 378.8, 370.8], 'attribute_3': ['Accept', 'Accept', 'Decline', 'Accept', 'Accept'], 'attribute_4': [False, False, True, False, True] }) # 构造待检查DataFrame check_df = pd.DataFrame({ 'my_id': ['ABC', 'ABS', 'ABB', 'ABB', 'APP'], 'parent_id': ['DEF', 'DES', 'DEG', 'DEG', 'DRE'], 'Attribute_1': ['A-', 'A-', 'A', 'C-', 'C-'], 'attribute2': [478.8, 388.8, 908.8, 378.8, 370.8], 'attr_3': ['Decline', 'Accept', 'Accept', 'Accept', 'Accept'], 'attribute_5': ['StRing', 'String', 'StrIng', 'String', 'STring'] }) # 定义列名映射:{基准列名: 待检查列名} col_mapping = { 'attribute_1': 'Attribute_1', 'attribute_2': 'attribute2', 'attribute_3': 'attr_3' } # 关联基准DF和待检查DF,只保留需要对比的基准列 merged_df = pd.merge( check_df, base_df[['my_id', 'parent_id'] + list(col_mapping.keys())], on=['my_id', 'parent_id'], how='left', suffixes=('', '_base') ) # 生成每个属性的错误标记列 for base_col, check_col in col_mapping.items(): merged_df[f'is_faulty_{check_col}'] = merged_df[check_col] != merged_df[f'{base_col}_base'] # 汇总错误属性列名 merged_df['faulty_attr'] = merged_df.apply( lambda row: [col for col in col_mapping.values() if row[f'is_faulty_{col}']], axis=1 ) # 筛选出有错误的行,保留需要的列 faulty_rows = merged_df[merged_df['faulty_attr'].apply(len) > 0][ ['my_id', 'parent_id'] + list(col_mapping.values()) + ['faulty_attr'] ] print(faulty_rows)
PySpark 实现方案
实现逻辑
- 定义列映射关系:同Pandas方案,明确基准与待检查列的对应关系。
- 关联两个DataFrame:使用
my_id和parent_id作为关联键进行左连接,将基准的对比属性带入待检查DF。 - 生成错误列名表达式:对每一组对应属性,用
when函数判断值是否不一致,不一致时返回待检查列名,一致时返回null。 - 汇总错误属性数组:将所有错误列名表达式组合成数组,用
array_remove移除空值,得到每行的错误属性列表。 - 筛选并整理结果:过滤出错误属性数组不为空的行,保留待检查DF的原始列和错误属性数组列。
代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, array, array_remove # 初始化SparkSession spark = SparkSession.builder.appName("FaultyRowCheck").getOrCreate() # 构造基准DataFrame base_data = [ ("ABC", "DEF", "A-", 378.8, "Accept", False), ("ABS", "DES", "A-", 388.8, "Accept", False), ("ABB", "DEG", "A", 908.8, "Decline", True), ("ABB", "DEG", "B-", 378.8, "Accept", False), ("APP", "DRE", "C-", 370.8, "Accept", True) ] base_df = spark.createDataFrame( base_data, ["my_id", "parent_id", "attribute_1", "attribute_2", "attribute_3", "attribute_4"] ) # 构造待检查DataFrame check_data = [ ("ABC", "DEF", "A-", 478.8, "Decline", "StRing"), ("ABS", "DES", "A-", 388.8, "Accept", "String"), ("ABB", "DEG", "A", 908.8, "Accept", "StrIng"), ("ABB", "DEG", "C-", 378.8, "Accept", "String"), ("APP", "DRE", "C-", 370.8, "Accept", "STring") ] check_df = spark.createDataFrame( check_data, ["my_id", "parent_id", "Attribute_1", "attribute2", "attr_3", "attribute_5"] ) # 定义列名映射:{基准列名: 待检查列名} col_mapping = { 'attribute_1': 'Attribute_1', 'attribute_2': 'attribute2', 'attribute_3': 'attr_3' } # 关联基准DF和待检查DF merged_df = check_df.join( base_df.select("my_id", "parent_id", *col_mapping.keys()), on=["my_id", "parent_id"], how="left" ) # 生成每个属性的错误列名表达式 faulty_exprs = [] for base_col, check_col in col_mapping.items(): faulty_exprs.append( when(col(check_col) != col(base_col), check_col).otherwise(None) ) # 汇总错误属性数组,移除空值 merged_df = merged_df.withColumn( "faulty_attr", array_remove(array(*faulty_exprs), None) ) # 筛选出有错误的行,保留需要的列 faulty_rows = merged_df.filter(col("faulty_attr").isNotNull() & (col("faulty_attr") != array())) \ .select("my_id", "parent_id", *col_mapping.values(), "faulty_attr") faulty_rows.show(truncate=False)
内容的提问来源于stack exchange,提问作者johnnydoe
相关产品推荐
相关产品推荐

