You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在PySpark中基于两个DataFrame实现数据更新与插入操作

PySpark 实现 DataFrame 增量更新(Upsert)方案

实现逻辑

你需要的是按col1作为主键的upsert操作,核心逻辑如下:

  • 对存量表df1和增量表df2做全外连接,关联键为col1
  • 按规则生成最终字段:
    • loaddate优先保留存量df1的取值,新增行取df2的loaddate
    • col2、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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 10:24:06