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

如何对比两个PySpark DataFrame以识别新增、更新与删除的记录

PySpark 更新记录查询实现方案

你可以通过以下两种常用方式实现需求,两种方式都能得到你期望的更新记录结果:

方法1:内连接+字段比对(灵活可控,支持自定义规则)

这种方法逻辑直观,还可以灵活调整比对逻辑,适合需要自定义差异判断规则的场景:

from pyspark.sql.functions import col

# 1. 给yesterday的非主键列加后缀,避免join后列名冲突
yesterday_suffix = yesterday.select(
    "item_id",
    *[col(c).alias(f"{c}_old") for c in yesterday.columns if c != "item_id"]
)

# 2. 按item_id内连接两天的数据
joined_df = today.join(yesterday_suffix, on="item_id", how="inner")

# 3. 筛选任意非主键列有差异的记录
compare_cols = [c for c in today.columns if c != "item_id"]
# 动态生成比对条件,不用手动写每一列的判断
filter_condition = None
for col_name in compare_cols:
    # 存在null值的场景可以用eqNullSafe替换!=,保证null值比对正确
    cond = col(col_name) != col(f"{col_name}_old")
    # cond = ~col(col_name).eqNullSafe(col(f"{col_name}_old"))
    filter_condition = cond if filter_condition is None else (filter_condition | cond)

# 4. 只保留今日的列输出
updated_records = joined_df.filter(filter_condition).select(today.columns)
updated_records.show()

方法2:exceptAll快速实现(简洁高效,适合无自定义规则场景)

如果只需要判断整行是否完全一致,不需要自定义比对规则,可以用更简洁的exceptAll算子实现:

# 先筛选出两天共有的item_id对应的记录
today_common = today.join(yesterday, on="item_id", how="left_semi")
yesterday_common = yesterday.join(today, on="item_id", how="left_semi")

# 取今天共有记录中 和 昨天共有记录 整行不一致的部分,即为更新记录
updated_records = today_common.exceptAll(yesterday_common)
updated_records.show()

两种方法运行后都会输出你期望的结果:

+-------+------+----+--------------+
|item_id|  name|cost|classification|
+-------+------+----+--------------+
|      1| Apple|5000|             A|
|      2|Banana|4000|             A|
|      3|Orange|3000|             B|
+-------+------+----+--------------+

注意事项

如果你的数据中存在null值,直接使用!=比对会出现判断错误,建议替换为~col(新列名).eqNullSafe(col(旧列名))的写法,可正确识别null值和非null值的差异。

内容的提问来源于stack exchange,提问作者Heber Brandao

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 12:06:03