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

如何从Databricks Delta表提取增量数据并基于版本实现增量备份?

利用Delta Lake变更数据馈送(CDF)实现增量数据加载与备份

核心原理

你已经为Delta表开启了delta.enableChangeDataFeed = true,这意味着Delta会自动记录每次版本变更的增量数据(包括插入、更新、删除操作)。完全可以利用这个特性实现增量备份,只需提取上一次运行到当前的版本/时间范围内的增量数据即可,无需同步全量。

增量数据加载到DataFrame的两种方式

1. 基于版本号提取

适合精准控制同步范围,需要记录上一次运行时的Delta表版本号(比如存在S3配置文件、数据库表中),然后从该版本到当前版本拉取增量。

2. 基于时间范围提取

适合定时任务场景,无需管理版本号,直接指定上一次运行时间到当前时间的范围即可。

操作步骤与示例代码

前提

确保Spark环境已集成Delta Lake,目标表已开启CDF。


示例1:基于版本号的增量加载与备份

from delta.tables import DeltaTable
import pyspark.sql.functions as F

# 1. 加载目标Delta表
delta_table = DeltaTable.forPath(spark, "/path/to/your/delta/table")

# 2. 获取上一次运行的版本号(示例值,实际从S3配置文件/数据库读取)
# 首次运行可设为0,每次同步完成后保存当前版本
last_sync_version = 3
current_version = delta_table.history(1).select("version").first()[0]

# 3. 加载增量数据到DataFrame
incremental_df = spark.read.format("delta") \
    .option("readChangeFeed", "true") \
    .option("startingVersion", last_sync_version) \
    .option("endingVersion", current_version) \
    .load("/path/to/your/delta/table")

# 可选:查看增量数据的操作类型(_change_type字段:insert/update_preimage/update_postimage/delete)
incremental_df.select("_change_type", "*").show()

# 4. 将增量数据写入S3备份(按版本分区,便于管理)
incremental_df.write.mode("append") \
    .partitionBy("_change_type") \
    .parquet(f"s3://your-backup-bucket/delta-increments/version_{last_sync_version}_to_{current_version}/")

# 5. 保存当前版本号到外部存储,供下次同步使用
spark.createDataFrame([(current_version,)], ["version"]) \
    .write.mode("overwrite").text("s3://your-config-bucket/last_sync_version.txt")

示例2:基于时间范围的增量加载与备份

from delta.tables import DeltaTable
import datetime

# 1. 定义时间范围(示例为1天前,实际从外部存储读取上次运行时间)
last_run_time = (datetime.datetime.now() - datetime.timedelta(days=1)).strftime("%Y-%m-%d %H:%M:%S")
current_run_time = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")

# 2. 加载增量数据
incremental_df = spark.read.format("delta") \
    .option("readChangeFeed", "true") \
    .option("startingTimestamp", last_run_time) \
    .option("endingTimestamp", current_run_time) \
    .load("/path/to/your/delta/table")

# 3. 写入S3备份(按时间分区)
incremental_df.write.mode("append") \
    .partitionBy("_change_type") \
    .parquet(f"s3://your-backup-bucket/delta-increments/time_{last_run_time.replace(' ', '_')}_to_{current_run_time.replace(' ', '_')}/")

# 4. 保存本次运行时间到外部存储
spark.createDataFrame([(current_run_time,)], ["run_time"]) \
    .write.mode("overwrite").text("s3://your-config-bucket/last_run_time.txt")

关键注意事项

  • 同步标记存储:必须可靠保存上一次的版本号或运行时间,避免重复同步或漏数据。推荐用数据库表(如MySQL)或S3文件存储,不要存在本地,防止节点故障丢失。
  • 变更类型过滤:CDF生成的update_preimage是更新前的旧数据,若只需要最终状态的增量,可过滤掉该类型:
    final_incremental_df = incremental_df.filter(F.col("_change_type").isin("insert", "update_postimage", "delete"))
    
  • S3写入优化:根据数据量选择Parquet/ORC格式,按日期或版本分区,方便后续恢复与查询。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 16:35:36