如何对比两个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
相关产品推荐
相关产品推荐

