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

如何使用Delta Lake Python库(非Spark)在S3上执行Upsert?

非Spark环境下Python实现Delta表Upsert的可行方案

Delta Lake原生的纯Python API(非PySpark)目前没有官方支持Merge/Upsert的直接接口,因为Delta的事务、元数据管理核心依赖Spark的分布式能力。不过你可以通过以下几种方案实现需求:

1. 基于delta-rs手动实现Upsert逻辑

delta-rs是Rust编写的Delta Lake库,提供Python绑定,支持在非Spark环境下读写S3上的Delta表。你可以手动实现Upsert流程:

  • 步骤:
    • 读取现有Delta表数据到Pandas DataFrame
    • 加载待更新的数据集
    • 基于主键过滤掉现有表中与新数据重复的记录,再合并新数据
    • 将合并后的数据写回Delta表(利用delta-rs的事务特性避免并发问题)
  • 示例代码:
import deltalake as dl
import pandas as pd

# 加载S3上的Delta表
delta_table = dl.DeltaTable("s3://your-bucket/path/to/delta-table")
existing_data = delta_table.to_pandas()

# 加载待Upsert的新数据
upsert_data = pd.read_csv("local-upsert-data.csv")

# 按主键(例如`id`)执行Upsert
# 先过滤掉现有数据中与新数据主键重复的行
filtered_existing = existing_data[~existing_data["id"].isin(upsert_data["id"])]
# 合并新数据
merged_data = pd.concat([filtered_existing, upsert_data], ignore_index=True)

# 写回Delta表,这里用overwrite模式,小数据量适用
delta_table.write(merged_data, mode="overwrite")
  • 注意:该方案适合数据量较小的场景,大数据量下手动处理会有性能瓶颈;同时要注意并发写入冲突,可依赖delta-rs的事务写入能力规避。

2. 借助Databricks REST API执行Merge

如果你的Delta表托管在Databricks上,可通过Databricks REST API提交Merge语句,由Databricks集群执行Upsert——无需自行维护Spark集群,按实际计算量付费,比分布式Spark处理小文件的成本更低。

3. 使用Local模式PySpark

若以上方案都不适用,可考虑用Local模式的PySpark在本地运行:

  • 优势:能直接使用Delta官方的Merge API,逻辑简洁,避免手动实现Upsert的潜在bug
  • 成本:本地运行无需分布式集群资源,处理小文件的开销远低于集群模式
  • 示例代码:
from pyspark.sql import SparkSession
from delta.tables import DeltaTable

# 初始化Local模式Spark会话
spark = SparkSession.builder \
    .appName("LocalDeltaUpsert") \
    .master("local[*]") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

# 加载S3上的Delta表
delta_table = DeltaTable.forPath(spark, "s3://your-bucket/path/to/delta-table")
# 加载待Upsert的数据
upsert_df = spark.read.csv("local-upsert-data.csv", header=True, inferSchema=True)

# 执行Merge/Upsert
delta_table.alias("target") \
    .merge(upsert_df.alias("source"), "target.id = source.id") \
    .whenMatchedUpdateAll() \
    .whenNotMatchedInsertAll() \
    .execute()

spark.stop()
  • 注意:适合数据量在本地硬件可承载的范围,比分布式集群更轻量化。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 18:45:21