基于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语句实现,步骤清晰:
- 基于水印列从Synapse视图读取增量数据;
- 将增量数据与目标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
相关产品推荐
相关产品推荐

