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

如何在Delta表中合并DataFrame实现插入、更新和删除操作?

Delta表插入/更新/删除合并方案

针对你的需求——以gene作为主键、disease为分区列,实现插入新记录、更新匹配记录、删除旧表存在但新DataFrame不存在的记录,以下是几种可行方案:

方案一:优化Merge语句实现全操作(推荐)

Delta Lake的MERGE命令支持删除操作,只需通过反向匹配标记待删除记录,一次操作完成所有逻辑:

  1. 以新DataFrame为源表,匹配旧Delta表的gene主键;
  2. 匹配到的记录执行更新;
  3. 源表中未匹配到的记录执行插入;
  4. 旧表中未匹配到的记录执行删除。

示例PySpark代码:

from delta.tables import DeltaTable

# 加载现有Delta表
delta_table = DeltaTable.forPath(spark, "/path/to/delta_table")
# 读取待合并的DataFrame
source_df = spark.read.csv("/path/to/source_data", header=True)

# 执行Merge全操作
delta_table.alias("target") \
    .merge(
        source_df.alias("source"),
        "target.gene = source.gene"
    ) \
    .whenMatchedUpdate(set={
        "value": "source.value",
        "disease": "source.disease"
    }) \
    .whenNotMatchedInsert(values={
        "disease": "source.disease",
        "gene": "source.gene",
        "value": "source.value"
    }) \
    .whenNotMatchedBySourceDelete()  # 核心:删除旧表中不在源DataFrame的记录
    .execute()

注意:whenNotMatchedBySourceDelete()是Delta Lake 2.0+版本支持的语法,若版本较低,可通过子查询匹配主键,手动删除不在源表的记录。

方案二:先删除再合并(需事务保证)

若你的Delta版本不支持上述语法,可采用「先删除旧表冗余记录,再执行Merge更新/插入」的方式,但必须在同一事务中执行,避免中间状态导致的数据不一致。

示例PySpark代码(同一作业内执行,默认原子性):

from delta.tables import DeltaTable

delta_table = DeltaTable.forPath(spark, "/path/to/delta_table")
source_df = spark.read.csv("/path/to/source_data", header=True)

# 提取源表的主键集合
source_genes = source_df.select("gene").distinct()

# 第一步:删除旧表中不在源表的记录
delta_table.delete("gene NOT IN (SELECT gene FROM source_genes)")

# 第二步:执行Merge完成更新/插入
delta_table.alias("target") \
    .merge(
        source_df.alias("source"),
        "target.gene = source.gene"
    ) \
    .whenMatchedUpdate(set={
        "value": "source.value",
        "disease": "source.disease"
    }) \
    .whenNotMatchedInsert(values={
        "disease": "source.disease",
        "gene": "source.gene",
        "value": "source.value"
    }) \
    .execute()

注:Spark作业内的连续Delta操作默认属于同一原子事务,不会暴露删除后、合并前的中间状态。

方案三:全量覆盖写入(适合小数据量)

如果数据量较小,可直接用新DataFrame覆盖写入Delta表,这种方式逻辑最简单,但仅适合无需保留旧表额外记录的场景:

source_df.write \
    .format("delta") \
    .mode("overwrite") \
    .option("mergeSchema", "true")  # 若有 schema 变更需开启
    .save("/path/to/delta_table")

缺点:大数据量下效率低,且无法保留操作历史痕迹。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 12:51:56