如何用Databricks PySpark增量读取ADLS中的Delta Parquet多文件
Delta 表增量读取实现方案
首次全量加载
首次加载时,直接读取ADLS路径下的Delta表并写入Spark表:
# 读取ADLS路径下的Delta表 full_df = spark.read.format("delta").load("/mnt/adls/path/to/delta-folder") # 写入Spark表(若表不存在则创建,存在则覆盖) full_df.write.mode("overwrite").saveAsTable("deltaTable.table")
增量读取新增/变更数据
Delta Lake的变更日志(Change Feed)支持增量读取新增或修改的数据,无需全量扫描。需确保Delta表已启用变更日志(Delta 2.0+默认自动启用,低版本需手动开启):
1. 启用变更日志(若未开启)
# 针对现有Delta表启用变更日志 spark.sql("ALTER TABLE deltaTable.table SET TBLPROPERTIES (delta.enableChangeFeed = true)") # 或在创建表时直接启用 spark.read.format("delta").load("/mnt/adls/path/to/delta-folder") \ .write.mode("overwrite") \ .option("delta.enableChangeFeed", "true") \ .saveAsTable("deltaTable.table")
2. 增量读取示例
使用readChangeFeed参数指定读取的版本范围或时间范围:
按版本号增量读取
# 获取当前表的最新版本 current_version = spark.sql("DESCRIBE HISTORY deltaTable.table").select("version").first()[0] # 读取上一次同步版本到当前版本的变更数据(假设上一次版本为last_processed_version) incremental_df = spark.read.format("delta") \ .option("readChangeFeed", "true") \ .option("startingVersion", last_processed_version) \ .option("endingVersion", current_version) \ .load("/mnt/adls/path/to/delta-folder") # 处理后更新last_processed_version值(可存储在数据库或Databricks作业参数中) last_processed_version = current_version
按时间范围增量读取
# 读取指定时间段内的变更数据(格式为yyyy-MM-dd HH:mm:ss) incremental_df = spark.read.format("delta") \ .option("readChangeFeed", "true") \ .option("startingTimestamp", "2024-01-01 00:00:00") \ .option("endingTimestamp", "2024-01-02 00:00:00") \ .load("/mnt/adls/path/to/delta-folder")
3. 合并增量数据到目标表
读取到增量数据后,可使用merge操作合并到已有的Spark表:
from delta.tables import DeltaTable target_table = DeltaTable.forName(spark, "deltaTable.table") target_table.alias("target").merge( incremental_df.alias("source"), "target.id = source.id" # 根据主键匹配 ).whenMatchedUpdateAll() \ .whenNotMatchedInsertAll() \ .execute()
注意事项
- 需记录每次增量读取的版本号或时间戳,避免重复处理数据
- 变更日志会记录所有数据操作(插入、更新、删除),可通过
_change_type字段区分操作类型(insert/update/delete)
内容的提问来源于stack exchange,提问作者Developer Rajinikanth
相关产品推荐
相关产品推荐

