使用PySpark对比两个DataFrame并输出不匹配字段的双边原值
PySpark 对比两个DataFrame获取不匹配行并排展示实现方案
需求说明
对比两个结构完全一致的DataFrame(df1、df2),以id为关联键,筛选出存在字段值不匹配的行,最终将不匹配字段的df1取值、df2取值并列展示。
输入样例
df1结构与数据:
+-----+---------+----------+--------+------+------+ | id|firstname|middlename|lastname|gender|salary| +-----+---------+----------+--------+------+------+ |42114| Robert| |Williams| M| 5000| |40288| Michael| Rose| | M| 4000| |39192| Maria| Anne| Jones| F| 4000| |36636| James| | Smith| M| 3000| | | Jen| Mary| Browln| F| -1| +-----+---------+----------+--------+------+------+
df2结构与数据:
+-----+---------+----------+--------+------+------+ | id|firstname|middlename|lastname|gender|salary| +-----+---------+----------+--------+------+------+ |42114| Robert| |Williams| M| 6000| |40288| Michael| Rose| | M| 4000| |39192| Maria| Anne| Jones| M| 4000| |36636| James| | Smith| M| 3000| | | Jen| Mary| Browln| F| -1| +-----+---------+----------+--------+------+------+
预期输出
仅展示存在不匹配的行,不匹配字段的双源值并列展示,匹配字段对应的衍生列为空:
+-----+---------+----------+--------+------+------+----------+----------+-----------+-----------+ | id|firstname|middlename|lastname|gender|salary|df1_gender|df2_gender|df1_salary |df2_salary | +-----+---------+----------+--------+------+------+----------+----------+-----------+-----------+ |42114| Robert| |Williams| M| 6000| null| null| 5000| 6000| |39192| Maria| Anne| Jones| M| 4000| F| M| null| null| +-----+---------+----------+--------+------+------+----------+----------+-----------+-----------+
原有代码问题
- select表达式嵌套列表,参数结构不符合Spark select方法的入参要求
- 仅生成了不匹配的列名数组,没有生成对应字段的双源值列
- 没有过滤完全匹配的行,会输出所有关联成功的行
修正后实现代码
from pyspark.sql import functions as F from pyspark.sql.functions import col, when, lit, array, array_remove # 提取排除关联键id之外的所有要对比的列 compare_cols = [c for c in df1.columns if c != 'id'] # 构造select字段列表 select_list = [col("id")] # 先保留df2的基础字段作为公共展示列(也可替换为df1的基础字段) for c in compare_cols: select_list.append(df2[c].alias(c)) # 为每个对比列生成df1、df2的取值列,仅不匹配时展示值,匹配时为null for c in compare_cols: select_list.append(when(df1[c] != df2[c], df1[c]).otherwise(lit(None)).alias(f"df1_{c}")) select_list.append(when(df1[c] != df2[c], df2[c]).otherwise(lit(None)).alias(f"df2_{c}")) # 生成不匹配列名数组,用于过滤完全匹配的行 mismatch_columns = array_remove( array(*[when(df1[c] != df2[c], lit(c)).otherwise(lit(None)) for c in compare_cols]), None ).alias("mismatch_columns") select_list.append(mismatch_columns) # 关联后过滤仅保留存在不匹配的行,输出结果 result_df = df1.join(df2, on="id", how="inner") \ .filter(F.size(mismatch_columns) > 0) \ .select(*select_list) result_df.show(truncate=False)
内容的提问来源于stack exchange,提问作者Mahima
相关产品推荐
相关产品推荐

