PySpark两个DataFrame按指定列匹配更新值的实现方案与效率咨询
PySpark 大小表合并最优方案
方案说明
你要实现的是典型的新表覆盖旧表匹配行、同时保留双方独有行的upsert需求,针对你data_old100亿行、data_new1000万行的大小表场景,过滤+Union方案远优于全连接方案,但不需要把匹配键收集到驱动端做列表检索,用Spark原生的左反连接(left anti join)实现分布式过滤即可,完全符合Spark设计理念。
最优实现代码:左反连接 + Union
实现逻辑
- 用左反连接过滤出
data_old中name、index没有在data_new中出现过的行,也就是仅在旧表存在的行 - 直接将过滤后的旧表和全量新表做Union,两个表Schema完全一致,不需要处理列冲突,得到的结果完全符合需求
from pyspark.sql import SparkSession from pyspark.sql.functions import col, broadcast spark = SparkSession.builder.appName("merge_df").getOrCreate() columns = ["name","index","name_2","index_2","string"] data_new = [(1,0,3,0,"C"),(1,1,3,1,"C"),(1,2,3,2,"C"),(1,3,3,3,"C"),(3,0,3,4,"C"),(3,1,3,5,"C")] data_old = [(1,0,1,0,"A"),(1,1,1,1,"A"),(1,2,1,2,"A"),(1,3,1,3,"A"),(2,0,2,0,"B"),(2,1,2,1,"B"),(2,2,2,2,"B"),(2,3,2,3,"B")] df_new = spark.createDataFrame(data_new, schema=columns) df_old = spark.createDataFrame(data_old, schema=columns) # 步骤1:左反连接过滤仅在旧表存在的行,主动广播小表优化,避免大表shuffle df_old_only = df_old.join(broadcast(df_new), on=["name", "index"], how="left_anti") # 步骤2:合并新表和过滤后的旧表 result = df_new.unionByName(df_old_only) # 验证结果 result.orderBy(col("name"), col("index")).show()
不推荐方案:全外连接
全外连接虽然也能实现需求,但需要对两个表全量数据做shuffle,100亿行的大表shuffle代价极高,仅作为参考实现:
from pyspark.sql.functions import coalesce df_full = df_old.alias("old").join(df_new.alias("new"), on=["name", "index"], how="full_outer") result_full = df_full.select( col("name"), col("index"), coalesce(col("new.name_2"), col("old.name_2")).alias("name_2"), coalesce(col("new.index_2"), col("old.index_2")).alias("index_2"), coalesce(col("new.string"), col("old.string")).alias("string") )
性能差异说明
- 左反连接+Union方案仅需要对1000万行的小表做shuffle,加上广播优化后100亿行的大表不需要任何shuffle,直接本地过滤,执行效率是全连接方案的数倍甚至数十倍
- 你最初想到的把匹配键收集为列表的方案会把全量匹配键拉到驱动端,容易触发驱动端OOM,同时用
isin过滤的性能远低于分布式的左反连接,不建议使用
内容的提问来源于stack exchange,提问作者Bondgirl
相关产品推荐
相关产品推荐

