PySpark:基于ID匹配用新DataFrame替换旧DataFrame对应行
解决PySpark数据框更新问题:替换匹配ID行并新增新行
数据说明
旧数据框 dataframe_old(简称dfo)
| Id | neighbor_sid | neighbor | division |
|---|---|---|---|
| a1 | 1100 | Naalehu | Hawaii |
| a2 | 1101 | key-west-fl | Miami |
| a3 | 1102 | lubbock | Texas |
| a10 | 1202 | bay-terraces | California |
新数据框 dataframe_new(简称dfn)
| Id | neighbor_sid | neighbor | division |
|---|---|---|---|
| a1 | 1100 | Naalehu | Hawaii |
| a2 | 1111 | key-largo-fl | Miami |
| a3 | 1103 | grapevine | Texas |
| a4 | 1115 | meriden-ct | Connecticut |
| a12 | 2002 | east-louisville | Kentucky |
需求目标
找出dfn中ID存在于dfo的行,用dfn的对应行替换dfo中的这些行;同时将dfn中dfo没有的ID的行新增进去,最终得到包含所有更新后行和原dfo中未被替换行的结果。
错误尝试代码
dfO.alias('a').join(dfN.alias('b'), on=['id'], how='left')\ .select( 'id', f.when( ~f.isnull(f.col('b.id')), f.col('b.id') ).otherwise(f.col('a.id')).alias('id'), 'b.col_3' )\ .union(dfO)\ .dropDuplicates()\ .sort('id')\ .show()
该代码逻辑混乱,仅针对单字段处理且未覆盖所有列,无法实现完整的行替换和新增需求。
正确解决方案
方法一:union + dropDuplicates(简单直观)
思路:将新数据框放在前面与旧数据框合并,再基于Id去重,保留先出现的新数据行。
from pyspark.sql import functions as f # 合并新、旧数据,新数据在前确保去重时被保留 combined_df = dfn.union(dfo) # 按Id去重并排序 final_df = combined_df.dropDuplicates(['Id']).sort('Id') final_df.show()
方法二:left_anti join筛选未更新行(性能更优)
思路:先提取旧数据中ID不在新数据里的行,再与新数据合并,避免全量数据union后去重的性能损耗。
from pyspark.sql import functions as f # 获取旧数据中未被新数据覆盖的行 old_unupdated = dfo.join(dfn, on=['Id'], how='left_anti') # 合并新数据和未更新的旧数据并排序 final_df = dfn.union(old_unupdated).sort('Id') final_df.show()
两种方法均可得到目标结果,方法二在数据量较大时性能更优,仅处理旧数据中未被替换的部分。
内容的提问来源于stack exchange,提问作者Ridhi
相关产品推荐
相关产品推荐

