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

如何比较两个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 22:33:15