Databricks中如何自动获取最新版本作为Change Feed的起止版本?
在Databricks中读取Delta Change Feed的最新单一版本
Delta Lake并没有提供直接的内置选项,能自动将最新的_commit_version同时设为startingVersion和endingVersion来读取单一版本的Change Feed,但你可以通过DeltaTable的元数据API手动获取最新版本号,再动态传入读取参数,完全不需要硬编码、水位表或流处理。
实现步骤:
- 通过
DeltaTable类获取目标表的元数据,提取最新版本号 - 将获取到的版本号动态传入
spark.read的参数中
完整代码示例:
from delta.tables import DeltaTable # 获取目标Delta表的最新版本号 delta_table = DeltaTable.forName(spark, "mdp_prd.bronze.nrq_customerassetproperty_autoloader_nodups") latest_version = delta_table.version() # 读取该最新版本的Change Feed数据 df1 = spark.read.format("delta") \ .option("readChangeFeed", "true") \ .option("startingVersion", latest_version) \ .option("endingVersion", latest_version) \ .table("mdp_prd.bronze.nrq_customerassetproperty_autoloader_nodups")
说明:
delta_table.version()会返回当前Delta表的最新提交版本号(long类型),可直接作为参数传入- 该方式为纯批处理操作,完全符合你不使用流处理的需求
- 无需维护额外的水位表,每次运行都会自动获取当前最新版本
内容的提问来源于stack exchange,提问作者Greencolor
相关产品推荐
相关产品推荐

