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

基于Delta Lake实现Azure Synapse到ADLS的增量数据加载方案咨询

针对Azure Synapse视图增量加载到ADLS分层架构的Delta Lake解决方案

是否应该使用Delta Lake?

完全适配你的需求,核心原因如下:

  • 完美支持增量加载+全量保留的Raw层诉求:Delta支持追加模式写入,同时自动维护全量数据,无需手动合并增量文件;
  • 原生提供**Upsert(Merge)**能力,直接解决Curated层的数据更新需求;
  • 可通过软删除标记或主键对比实现源端删除记录的同步;
  • 具备ACID事务、时间旅行、小文件合并等生产级特性,适配分层数据架构的稳定性要求。

如何向Delta Lake执行Upsert操作?

Upsert通过Spark SQL的MERGE INTO语句实现,步骤清晰:

  1. 基于水印列从Synapse视图读取增量数据;
  2. 将增量数据与目标Delta表按业务主键匹配,执行"匹配则更新/删除,不匹配则插入"逻辑。

代码示例(PySpark)

# 1. 获取上次同步的水印值(可从Delta事务日志或外部元数据表读取)
last_watermark = "2024-05-01 00:00:00"

# 2. 读取Synapse视图的增量数据
incremental_df = spark.read \
    .format("com.microsoft.sqlserver.jdbc.spark") \
    .option("url", "jdbc:sqlserver://<synapse-workspace>.sql.azuresynapse.net:1433;database=<db-name>") \
    .option("dbtable", "(SELECT * FROM dbo.your_view WHERE last_modified_time > '{}') AS inc_data".format(last_watermark)) \
    .option("user", "<username>") \
    .option("password", "<password>") \
    .load()

# 3. 注册临时视图用于Merge操作
incremental_df.createOrReplaceTempView("incremental_data")

# 4. 执行Upsert到Curated层Delta表
spark.sql("""
    MERGE INTO curated.your_curated_table t
    USING incremental_data s
    ON t.id = s.id -- 替换为你的业务主键
    WHEN MATCHED AND s.is_deleted = 1 THEN DELETE -- 源端有软删除标记时直接删除
    WHEN MATCHED THEN UPDATE SET * -- 匹配则全量更新(也可指定具体列)
    WHEN NOT MATCHED THEN INSERT * -- 不匹配则插入新记录
""")

# 5. 更新水印值到元数据存储,供下次同步使用

如何同步源端删除的记录?

分两种场景处理:

场景1:源端支持软删除(推荐)

如果Synapse视图包含is_deleted或类似标记列,直接在上述MERGE INTO语句中加入WHEN MATCHED AND s.is_deleted = 1 THEN DELETE即可,这是最高效的同步方式。

场景2:源端为硬删除

若源端直接删除记录且无标记,需定期对比源端与Delta表的主键集合,找出缺失的主键并删除:

# 读取源端所有业务主键
source_keys = spark.read \
    .format("com.microsoft.sqlserver.jdbc.spark") \
    .option("url", "<synapse-jdbc-url>") \
    .option("dbtable", "SELECT id FROM dbo.your_view") \
    .option("user", "<username>") \
    .option("password", "<password>") \
    .load() \
    .select("id")

# 读取Delta表所有业务主键
delta_keys = spark.read \
    .format("delta") \
    .load("/mnt/adls/curated/your_curated_table") \
    .select("id")

# 找出Delta表中存在但源端已删除的主键
deleted_keys = delta_keys.join(source_keys, delta_keys.id == source_keys.id, "left_anti")

# 删除Delta表中的对应记录
deleted_keys.createOrReplaceTempView("deleted_keys")
spark.sql("DELETE FROM curated.your_curated_table WHERE id IN (SELECT id FROM deleted_keys)")

注意:硬删除同步建议按天或固定周期执行,避免频繁全量扫描影响性能。

合适的分区列选择

分区列的核心原则是匹配数据加载和查询的模式,具体建议:

Raw层分区

  • 优先选择水印列的日期维度(比如date(last_modified_time)):因为增量加载按时间触发,分区后每次仅写入对应日期的分区,后续全量导出也可按分区并行处理;
  • 避免选择高基数列(如用户ID、订单ID),会导致小文件过多。

Curated层分区

  • 优先选择时间维度+高频查询的业务维度:比如如果Curated层经常按日期和业务线查询,就按modified_date和business_line分区;
  • 控制分区粒度:每个分区的数据量建议在1-10GB之间,避免分区过多或过少;
  • 若无明显业务查询维度,仅按日期分区即可。

Raw层全量数据的处理

Raw层需保留全量数据且通过追加更新,直接用Delta表的append模式写入即可:

# 追加增量数据到Raw层Delta表
incremental_df.write \
    .format("delta") \
    .mode("append") \
    .partitionBy("modified_date") \
    .save("/mnt/adls/raw/your_raw_table")

# 如需导出全量文件到ADLS,直接读取整个Delta表写入即可
spark.read.format("delta").load("/mnt/adls/raw/your_raw_table") \
    .write \
    .format("parquet") \
    .mode("overwrite") \
    .save("/mnt/adls/raw/full_export")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 10:20:41