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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 23:30:31