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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 05:36:35