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

如何通过Spark DataFrame实现数据库表的行插入与更新操作

Spark 实现Upsert(存在更新、不存在插入)的落地方法

你提到的两种思路都可落地,分别对应不同的适用场景,具体实现方法如下:

方案1:全量合并重写

适合目标表数据量不大(一般10G以内)、或者是小维度表的场景,实现成本最低,不依赖数据库特性。

  • 第一步:读取目标库的原表数据作为DataFrame,记为old_df
  • 第二步:将新数据new_df和old_df基于键列Key1、Key2做全外连接
  • 第三步:合并值字段,优先取新数据的Value,无新数据时保留原表Value,PySpark示例代码:
from pyspark.sql.functions import coalesce

merged_df = new_df.join(old_df, on=["Key1", "Key2"], how="outer") \
                  .select(
                      "Key1",
                      "Key2",
                      coalesce(new_df.Value, old_df.Value).alias("Value")
                  )
  • 第四步:将合并后的merged_df用overwrite模式覆写目标表即可。

该方案的优势是实现简单,兼容MySQL、PostgreSQL、Hive等所有存储介质;缺点是全量读、全量写的性能开销大,不适用于大表同步。

方案2:增量执行插入/更新

适合目标表数据量大的场景,无需读取全量表,仅处理增量数据,依赖目标存储的Upsert特性。

场景A:写入支持原生Upsert语法的关系型数据库(MySQL、PostgreSQL等)

实现逻辑是通过Spark的foreachPartition算子批量处理分区数据,拼接对应数据库的Upsert语句执行:

注意:目标表必须提前给Key1、Key2创建联合唯一主键,否则无法触发更新逻辑。

以MySQL为例示例代码:

import pymysql

def upsert_partition(rows):
    # 每个分区创建一次数据库连接,避免频繁创建连接消耗性能
    conn = pymysql.connect(host="数据库地址", user="用户名", password="密码", database="库名")
    cursor = conn.cursor()
    sql = """
    INSERT INTO target_table (Key1, Key2, Value) 
    VALUES (%s, %s, %s)
    ON DUPLICATE KEY UPDATE Value = VALUES(Value)
    """
    batch_data = [(row.Key1, row.Key2, row.Value) for row in rows]
    cursor.executemany(sql, batch_data)
    conn.commit()
    cursor.close()
    conn.close()

new_df.foreachPartition(upsert_partition)

PostgreSQL只需要将SQL替换为对应语法即可:INSERT INTO ... ON CONFLICT (Key1,Key2) DO UPDATE SET Value = EXCLUDED.Value。

场景B:写入数据湖表(Iceberg、Hudi、Delta Lake)

三种主流数据湖格式都原生支持Spark的Upsert操作,以Delta Lake为例:

from delta.tables import DeltaTable

# 读取目标Delta表
delta_table = DeltaTable.forPath(spark, "目标表存储路径")

# 执行merge操作
delta_table.alias("old") \
  .merge(
    new_df.alias("new"),
    "old.Key1 = new.Key1 AND old.Key2 = new.Key2"
  ) \
  .whenMatchedUpdate(set = {"Value": "new.Value"}) \
  .whenNotMatchedInsert(values = {
    "Key1": "new.Key1",
    "Key2": "new.Key2",
    "Value": "new.Value"
  }) \
  .execute()

该方案性能最优,适合大数据量的离线、实时同步场景。

选型参考

  • 小表/维度表同步优先选方案1,实现简单不易出错
  • 大表同步至关系型数据库选方案2的foreachPartition拼接Upsert语句的方式
  • 大表同步至数据湖存储,直接使用对应数据湖的merge语法即可

内容的提问来源于stack exchange,提问作者Alexander Lopatin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 11:06:03