批处理模式下如何让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
相关产品推荐
相关产品推荐

