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

如何在Azure环境中实现Bronze到Silver层的增量数据UPSERT加载?

Azure环境下Bronze到Silver层增量UPSERT实现方案

前置条件

  • Silver层已创建Delta表,需包含主键/唯一键以及修改日期字段,用于匹配增量数据
  • 控制表需记录上次同步的最大修改日期(比如字段last_sync_max_modify_date),作为增量提取的基准

步骤1:增量数据提取与Bronze层落地

  • 基于控制表的last_sync_max_modify_date,从源系统拉取修改日期大于该值的增量数据
  • 将增量数据以Parquet格式写入ADLS Gen2的Bronze层,建议按日期分区(比如/bronze/table_name/year=YYYY/month=MM/day=DD/),后续读取更高效

步骤2:Silver层Delta表的UPSERT操作

Azure里常用的Delta Lake处理平台是Databricks和Synapse Analytics,以下是两种场景的实现方式:

场景1:用Azure Databricks实现

  • 读取Bronze层的增量Parquet数据:
from pyspark.sql.functions import col

# 替换为你的ADLS Gen2实际路径
bronze_incremental_df = spark.read.parquet("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/bronze/table_name/[分区路径]")
  • 执行Delta表的Merge(即UPSERT)操作:
# 注册临时视图,用SQL编写Merge逻辑更直观
bronze_incremental_df.createOrReplaceTempView("bronze_incremental")

# 匹配主键后更新现有记录,不匹配则插入新记录
spark.sql("""
    MERGE INTO silver_table st
    USING bronze_incremental bi
    ON st.primary_key = bi.primary_key
    WHEN MATCHED THEN UPDATE SET *
    WHEN NOT MATCHED THEN INSERT *
""")

提示:如果只需要更新特定字段,可把UPDATE SET *改成UPDATE SET st.field1 = bi.field1, st.field2 = bi.field2,避免覆盖不需要更新的字段

场景2:用Azure Synapse Analytics实现

  • 先确保Synapse已配置ADLS Gen2链接服务,且Spark池启用Delta Lake支持
  • 读取增量数据并执行Merge:
# 读取Bronze层增量Parquet文件
bronze_df = spark.read.parquet("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/bronze/table_name/")

# 关联Silver层Delta表执行Merge操作
from delta.tables import DeltaTable

silver_delta_table = DeltaTable.forPath(spark, "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/silver/table_name")

silver_delta_table.alias("st").merge(
    bronze_df.alias("bi"),
    "st.primary_key = bi.primary_key"
).whenMatchedUpdateAll().whenNotMatchedInsertAll().execute()

步骤3:更新控制表

  • 从本次增量数据中提取最大的修改日期:
max_modify_date = bronze_incremental_df.agg({"modify_date": "max"}).collect()[0][0]
  • 将该日期更新到控制表,作为下次同步的基准,保证增量数据不重复、不遗漏

优化建议

  • Bronze层Parquet按修改日期分区,减少每次读取的数据量
  • 定期对Silver层Delta表做优化和清理,提升性能:
OPTIMIZE silver_table ZORDER BY primary_key;
VACUUM silver_table RETAIN 7 DAYS;
  • 用Azure Data Factory (ADF) 调度整个流程,实现自动化增量同步

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 07:01:03