如何在PySpark中对比两个DataFrame以获取更新及新增记录
实现方案
核心逻辑是以NAME为主键,优先保留增量表df3的记录,存量表df2仅保留增量表中不存在的主键对应记录,下面提供两种常用实现方式:
方法1:过滤存量 + 增量合并(性能更优,适合增量数据量远小于存量的场景)
逻辑:先把存量df2里不需要更新的记录(也就是NAME不在df3里的记录)筛出来,再和全量增量df3合并即可。
PySpark 代码
# 先取df3的所有NAME集合 exist_name_df = df3.select("NAME") # 筛出df2中NAME不在df3里的记录 df2_remain = df2.join(exist_name_df, on="NAME", how="left_anti") # 合并保留的存量记录和全量增量记录 result_df = df2_remain.unionByName(df3) # 验证结果 result_df.show()
Spark SQL 代码
SELECT * FROM df2 WHERE NAME NOT IN (SELECT NAME FROM df3) UNION ALL SELECT * FROM df3
执行后输出就是你需要的结果:
+----+-------+------+ |NAME|BALANCE|SALARY| +----+-------+------+ |Liza| 20| 900| |PPan| 10| 700| | Cal| 70| 888| +----+-------+------+
方法2:全外连接 + 字段优先级取值(适合字段多、后续字段变更频繁的场景)
逻辑:两个表全外连接后,非主键字段优先取df3的值,df3为空时取df2的原值。
PySpark 代码
from pyspark.sql.functions import coalesce # 全外连接两个表 joined_df = df2.join(df3, on="NAME", how="outer") # 按优先级合并字段 result_df = joined_df.select( "NAME", coalesce(df3.BALANCE, df2.BALANCE).alias("BALANCE"), coalesce(df3.SALARY, df2.SALARY).alias("SALARY") ) # 验证结果 result_df.show()
两种方法输出结果完全一致,你可以根据自己的实际数据量和表结构灵活选择。
内容的提问来源于stack exchange,提问作者greenking
相关产品推荐
相关产品推荐

