Synapse Delta表:Silver到Gold加载时如何读取每条记录最新版本?
Synapse Delta表读取每个键的最新版本记录及版本号获取方案
一、读取每个键(Key)的最新版本记录
在Silver到Gold层的ETL流程中,利用Delta Lake的特性可高效筛选每个键的最新记录,以下是两种常用方案:
1. 窗口函数筛选最新记录
通过窗口函数按键分区、按时间戳降序排序,筛选出每个分区的第一条记录:
from pyspark.sql.window import Window from pyspark.sql.functions import row_number, col # 读取Silver层Delta表 silver_df = spark.read.format("delta").load("abfss://<容器名>@<存储账户>.dfs.core.windows.net/silver/<表名>") # 定义窗口规则:按业务键分区,按更新时间戳降序排列 window_spec = Window.partitionBy("your_key_column").orderBy(col("update_timestamp").desc()) # 筛选每个键的最新记录 gold_df = silver_df.withColumn("row_rank", row_number().over(window_spec)) \ .filter(col("row_rank") == 1) \ .drop("row_rank") # 写入Gold层Delta表 gold_df.write.format("delta").mode("overwrite").save("abfss://<容器名>@<存储账户>.dfs.core.windows.net/gold/<表名>")
说明:
update_timestamp可替换为Delta表内置的_commit_timestamp(基于数据提交时间判断最新),或业务自定义的更新时间字段。
2. 使用Merge语句增量同步最新记录
如果需要增量更新Gold层而非全量覆盖,推荐使用Delta的Merge语句,仅更新/插入比Gold层更晚的记录:
# 定义Gold层表路径 gold_table_path = "abfss://<容器名>@<存储账户>.dfs.core.windows.net/gold/<表名>" # 读取Silver层数据 silver_df = spark.read.format("delta").load("abfss://<容器名>@<存储账户>.dfs.core.windows.net/silver/<表名>") # 创建临时视图用于SQL Merge silver_df.createOrReplaceTempView("silver_temp") # 执行Merge操作 spark.sql(f""" MERGE INTO delta.`{gold_table_path}` gold USING silver_temp silver ON gold.your_key_column = silver.your_key_column WHEN MATCHED AND silver.update_timestamp > gold.update_timestamp THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT * """)
二、获取Delta表的最新版本号(替代Databricks方案)
针对Synapse PySpark环境,以下两种方式可获取Delta表的最新版本号:
1. 利用DeltaTable API查询历史
from delta.tables import DeltaTable # 加载目标Delta表 delta_table = DeltaTable.forPath(spark, "abfss://<容器名>@<存储账户>.dfs.core.windows.net/<表路径>") # 获取表的版本历史,筛选最新版本 history_df = delta_table.history() latest_version = history_df.orderBy(col("version").desc()).select("version").first()[0] print(f"最新版本号:{latest_version}")
2. Spark SQL查询元数据
直接通过SQL查询Delta表的历史记录并聚合得到最新版本:
-- 查询全量历史 DESCRIBE HISTORY delta.`abfss://<容器名>@<存储账户>.dfs.core.windows.net/<表路径>` -- 直接获取最新版本号 SELECT MAX(version) AS latest_version FROM (DESCRIBE HISTORY delta.`abfss://<容器名>@<存储账户>.dfs.core.windows.net/<表路径>`)
内容的提问来源于stack exchange,提问作者david
相关产品推荐
相关产品推荐

