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

SQL Server数据导入数据湖后,如何实现数据合并(Upsert)?

数据湖Parquet文件的Upsert(合并新增/更新数据)方案

核心逻辑

Upsert本质是插入新数据+更新已有匹配数据,但纯Parquet是不可变文件,无法直接修改,通常通过「转换为支持事务的格式」或「合并后重写文件」的方式实现,结合增量同步逻辑优化效率。

具体实现方案

1. 使用Delta Lake(推荐方案)

Delta Lake基于Parquet格式扩展,支持ACID事务和原生Upsert操作,是数据湖合并场景的最优选择:

  • 先将现有Parquet数据转换为Delta表:
    from delta.tables import DeltaTable
    from pyspark.sql import SparkSession
    
    spark = SparkSession.builder.appName("ParquetUpsert").getOrCreate()
    
    # 读取已有Parquet文件,转存为Delta表
    parquet_data = spark.read.parquet("/data-lake/your-parquet-folder")
    parquet_data.write.format("delta").mode("overwrite").save("/data-lake/your-delta-table")
    
  • 执行Upsert操作:先从SQL Server拉取增量数据(通过时间戳、CDC或自增ID过滤),再与Delta表合并:
    # 拉取SQL Server增量数据(示例用last_updated_time过滤)
    incremental_data = spark.read.jdbc(
        url="jdbc:sqlserver://your-sql-server:1433;databaseName=your-db",
        table="(SELECT * FROM your_source_table WHERE last_updated_time > '上次同步的时间戳') AS inc_data",
        properties={"user": "your-username", "password": "your-password"}
    )
    
    # 加载Delta表并执行Upsert
    delta_table = DeltaTable.forPath(spark, "/data-lake/your-delta-table")
    delta_table.alias("target").merge(
        incremental_data.alias("source"),
        "target.id = source.id"  # 替换为你的唯一主键字段
    ).whenMatchedUpdateAll(  # 匹配到则更新所有字段
    ).whenNotMatchedInsertAll(  # 未匹配到则插入
    ).execute()
    
  • 后续可直接用Delta表查询,也可按需导出为Parquet格式。

2. 纯Parquet文件的Upsert(无Delta Lake)

如果无法引入Delta Lake,只能通过「合并重写」实现,适合数据量较小的场景:

  • 步骤1:读取数据湖全量Parquet数据 + SQL Server增量数据
  • 步骤2:按主键去重,保留最新版本(依赖更新时间戳字段):
    # 读取全量Parquet数据
    full_parquet_data = spark.read.parquet("/data-lake/your-parquet-folder")
    # 读取增量数据(同上面的JDBC逻辑)
    incremental_data = spark.read.jdbc(...)
    
    # 合并数据集,按主键分组保留最新记录
    combined_data = full_parquet_data.unionByName(incremental_data)
    final_data = combined_data.groupBy("id").agg(
        *[max(c).alias(c) if c == "last_updated_time" else first(c).alias(c) for c in combined_data.columns]
    )
    
  • 步骤3:将最终数据重写回数据湖(会覆盖原有文件):
    final_data.write.mode("overwrite").parquet("/data-lake/your-parquet-folder")
    
  • 优化:如果数据按分区存储(比如按日期),可只重写有增量的分区,减少IO开销。

3. 增量同步的前置优化

  • 开启SQL Server的CDC(变更数据捕获),精准获取新增、更新、删除的变更记录,避免全量扫描源表。
  • 在数据管道中记录每次同步的水印值(比如最大的last_updated_time或CDC的LSN号),下次同步仅拉取水印之后的数据。

注意事项

  • 必须有唯一主键(如ID),否则无法准确匹配需要更新的记录。
  • 纯Parquet重写方式在数据量大时性能较差,优先考虑Delta Lake。
  • 若需处理删除操作:Delta Lake可在merge中添加whenMatchedDelete条件;纯Parquet则需在合并后过滤已删除记录再重写。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 21:40:40