如何在PySpark中基于两个DataFrame实现数据更新与插入操作
PySpark 实现 DataFrame 增量更新(Upsert)方案
实现逻辑
你需要的是按col1作为主键的upsert操作,核心逻辑如下:
- 对存量表df1和增量表df2做全外连接,关联键为
col1 - 按规则生成最终字段:
loaddate优先保留存量df1的取值,新增行取df2的loaddatecol2、lastupdatedate优先取增量df2的新值,存量未变更行保留df1原值
- 筛选出最终字段即可得到目标结果
代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when # 初始化SparkSession,实际使用时可替换为你自己的Spark上下文 spark = SparkSession.builder.appName("data_upsert").getOrCreate() # -------------------------- # 以下为示例构造数据,实际使用时替换为你自己的df1、df2读取逻辑 # df1 = spark.read.parquet("你的存量表路径") # df2 = spark.read.parquet("你的增量表路径") # -------------------------- df1_data = [ ("a", 1, "12/02/21", "12/02/21"), ("b", 2, "12/02/21", "12/02/21"), ("c", 3, "12/02/21", "12/02/21"), ("d", 4, "12/02/21", "12/02/21"), ("e", 5, "12/02/21", "12/02/21") ] df1 = spark.createDataFrame(df1_data, schema=["col1", "col2", "loaddate", "lastupdatedate"]) df2_data = [ ("a", 10, "12/12/21", "12/12/21"), ("f", 2, "12/12/21", "12/12/21"), ("g", 3, "12/12/21", "12/12/21") ] df2 = spark.createDataFrame(df2_data, schema=["col1", "col2", "loaddate", "lastupdatedate"]) # 核心upsert逻辑 joined_df = df1.alias("old").join(df2.alias("new"), on="col1", how="outer") result_df = joined_df.select( col("col1"), when(col("new.col2").isNotNull(), col("new.col2")).otherwise(col("old.col2")).alias("col2"), when(col("old.loaddate").isNotNull(), col("old.loaddate")).otherwise(col("new.loaddate")).alias("loaddate"), when(col("new.lastupdatedate").isNotNull(), col("new.lastupdatedate")).otherwise(col("old.lastupdatedate")).alias("lastupdatedate") ) # 验证结果 result_df.orderBy("col1").show()
生产环境优化建议
如果是大数据量的生产场景,推荐使用Delta Lake的mergeInto语法实现,性能更高,还支持事务保证:
# 假设你已经将存量数据保存为Delta表 from delta.tables import DeltaTable delta_table = DeltaTable.forPath(spark, "你的存量Delta表路径") delta_table.alias("old") \ .merge( df2.alias("new"), "old.col1 = new.col1" ) \ .whenMatchedUpdate(set = { "col2": "new.col2", "lastupdatedate": "new.lastupdatedate" }) \ .whenNotMatchedInsertAll() \ .execute()
内容的提问来源于stack exchange,提问作者Ratnakar Mogulla
相关产品推荐
相关产品推荐

