如何比较两个PySpark DataFrame并合并双方的不匹配数据?
解决Spark DataFrame双向不匹配记录合并问题
我明白你的需求啦——你现在用subtract只能拿到df2里存在但df1没有的记录,但其实你需要汇总df1和df2两边所有不匹配的行,也就是df1独有的、df2独有的记录都要整合到一起,对吧?
先说说当前代码的问题
df2.subtract(df1)这个操作只会返回df2中不存在于df1的整行记录,反过来df1里独有的记录就被漏掉了,所以只能拿到单向的差异。
解决方案:双向差异+合并
要拿到完整的双向不匹配,我们需要分别计算两边的独有记录,再把它们合并起来,还可以给每条记录加个来源标识,方便区分是来自哪一边的差异。
完整代码示例
from pyspark.sql.functions import lit # 读取数据 df1 = spark.read.load("df1.csv", format="csv", header="True") df2 = spark.read.load("df2.csv", format="csv", header="True") # 获取df1独有的记录,并添加来源标记 df1_unique = df1.subtract(df2).withColumn("record_source", lit("df1")) # 获取df2独有的记录,并添加来源标记 df2_unique = df2.subtract(df1).withColumn("record_source", lit("df2")) # 合并两个差异结果:推荐用unionByName,避免列顺序不一致导致的错误 all_mismatches = df1_unique.unionByName(df2_unique) # 查看最终合并的不匹配记录 all_mismatches.show(truncate=False)
关键细节说明
subtract操作是整行去重比较,只有当两行所有字段的值都完全一致时,才会被判定为匹配。unionByName比普通union更安全:如果df1和df2的列顺序不一样,union会按列的位置匹配,容易出现数据错位;unionByName会严格按照列名匹配,更可靠。- 添加
record_source列是可选但很实用的操作,能让你一眼看出这条差异记录原本属于哪个DataFrame。
补充:如果需要字段级别的不匹配(而非整行)
如果你的需求是找出同一主键下字段值不一致的行(比如同一条id的记录,某个字段值在df1和df2里不一样),那可以用全外连接的方式:
# 假设你的主键列是`id`,需要比较的字段是`col1`、`col2` joined_df = df1.join(df2, on="id", how="full_outer", suffixes=("_df1", "_df2")) # 过滤出字段不匹配或只存在于一边的记录 field_mismatches = joined_df.filter( # 字段值不相等的情况 (joined_df["col1_df1"] != joined_df["col1_df2"]) | (joined_df["col2_df1"] != joined_df["col2_df2"]) | # 只存在于df1或df2的情况 joined_df["id"].isNull() | joined_df["col1_df1"].isNull() | joined_df["col1_df2"].isNull() ) field_mismatches.show(truncate=False)
内容的提问来源于stack exchange,提问作者Suneel Kumar
相关产品推荐
相关产品推荐

