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

批处理模式下如何让Azure Databricks Spark仅读取Azure Data Lake Gen2新增数据

解决方案

这里提供3种生产环境常用的批处理增量读取方案,均不需要引入流式处理逻辑:

方案1:使用Databricks Auto Loader 增量批处理模式(最推荐)

Auto Loader 是Databricks官方提供的文件源工具,默认支持增量追踪,不需要自行维护元数据,也可以作为单次批处理任务运行,无需常驻流进程:

  • 核心逻辑:Auto Loader会自动在你指定的检查点目录中记录所有已处理的文件信息,每次任务运行时仅读取上次运行结束后新增的文件
  • 代码示例:
# 读取新增文件,将格式、路径替换为实际业务值
df = spark.read.format("cloudFiles") \
  .option("cloudFiles.format", "parquet") \
  .option("cloudFiles.schemaLocation", "/checkpoint/schema/") \
  .option("cloudFiles.includeExistingFiles", "false") # 首次运行是否处理存量文件可按需调整
  .load("/your/adls/storage/path/")

# 业务处理逻辑
result_df = df.transform(your_business_logic)

# 写出结果
result_df.write.format("delta").mode("append").save("/output/path/")

# 任务结束后Auto Loader会自动更新检查点的已处理文件记录,下次运行自动过滤
  • 优点:无需自行维护增量元数据,支持大数量级文件,自动处理Schema演变,官方原生支持维护成本低
  • 注意:不要修改或删除检查点目录,否则会触发全量文件读取

方案2:基于时间水位线自行控制增量

如果你的数据湖文件按时间分区存储(比如目录结构为/dt=2024-05-01/)或者文件有明确的修改时间标记,可以自己维护增量水位线:

  • 核心逻辑:将每次处理的最大时间(分区时间/文件修改时间)存储在独立的控制表/配置文件中,下次任务运行时先读取该水位线,仅拉取时间大于水位线的文件
  • 代码示例:
from pyspark.sql.functions import col, max

# 1. 读取上次处理的水位线,示例为存在Delta控制表的场景
last_process_time = spark.read.table("process_control_table").select("max_dt").collect()[0][0]

# 2. 只读取大于水位线的分区数据
df = spark.read.parquet("/your/data/path/") \
  .filter(col("dt") > last_process_time)

# 3. 业务处理
result_df = df.transform(your_business_logic)

# 4. 更新水位线到控制表
new_max_dt = df.select(max("dt")).collect()[0][0]
spark.createDataFrame([(new_max_dt,)], ["max_dt"]) \
  .write.format("delta").mode("overwrite").saveAsTable("process_control_table")
  • 优点:逻辑简单透明,自定义程度高,不依赖Databricks专属功能
  • 注意:如果存在晚到的小于当前水位线的数据,需要单独开发补数据逻辑

方案3:基于已处理文件列表差集过滤

适合文件只会新增、不会修改/删除的场景:

  • 核心逻辑:将所有已经处理过的文件路径存储在专属元数据表中,每次运行时先拉取数据湖当前所有文件路径,和元数据表做差集得到新增文件列表,仅读取这些文件,处理完成后把新的文件路径追加到元数据表
  • 代码示例:
from pyspark.sql.functions import input_file_name

# 1. 读取已处理的文件列表
processed_files = spark.read.table("processed_files_table").select("file_path")

# 2. 列出当前数据湖所有文件,过滤出未处理的
all_files = spark.read.parquet("/your/data/path/").select(input_file_name().alias("file_path")).distinct()
new_files = all_files.join(processed_files, on="file_path", how="left_anti")
new_file_list = [row.file_path for row in new_files.collect()]

# 3. 只读取新增文件
df = spark.read.parquet(*new_file_list)

# 4. 业务处理
result_df = df.transform(your_business_logic)

# 5. 追加新处理的文件路径到元数据表
new_files.write.format("delta").mode("append").saveAsTable("processed_files_table")
  • 优点:逻辑完全可控,适合文件路径有明确规则的场景
  • 注意:文件数量特别大时,拉取全量文件列表会有额外性能损耗

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 10:48:00