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

如何使用Databricks Autoloader加载特定前缀的CSV文件

解决Databricks Autoloader加载子目录下特定前缀CSV文件的问题

问题场景

使用Databricks Autoloader加载Azure Blob容器中/Dir1/Dir2目录下(实际文件存储于其子目录)、文件名以String_Pattern_开头的CSV文件时遇到以下问题:

  • 配置pathGlobfilter选项设置过滤规则后,始终无法加载到目标文件;
  • 移除过滤规则或仅过滤*.csv时,可正常加载所有文件;
  • 尝试用/**/String_Pattern_*.csv作为过滤规则适配子目录结构,收到报错:

details = "Cannot infer schema when the input path abfss://<container-path>/Dir1/Dir2/ is empty. Please try to start the stream when there are files in the input path, or specify the schema."

原代码如下:

df = (spark.readStream
  .format("cloudFiles") 
  .option("cloudFiles.format", "csv")
  .option("header", "true") 
  .option("cloudFiles.includeExistingFiles", "true") 
  .option("cloudFiles.inferColumnTypes", "true") 
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .option("pathGlobfilter", "String_Pattern_*.csv")
  .load("abfss://<container-path>/Dir1/Dir2/")
)

原因分析

  1. 参数名错误:使用了Spark原生的pathGlobfilter(小写f)而非Autoloader专用的cloudFiles.pathGlobFilter参数,导致过滤规则未被正确识别;
  2. 过滤规则误用:pathGlobFilter仅用于匹配文件名,不支持路径通配符(如/**/),错误的规则会导致没有匹配到任何文件;
  3. Schema推断失败:当没有匹配到文件时,Autoloader无法自动推断Schema,从而抛出路径为空的错误。

解决方案

方案1:使用Autoloader专用过滤参数

指定正确的cloudFiles.pathGlobFilter参数,并确保开启递归扫描子目录(默认已开启,可明确配置):

df = (spark.readStream
  .format("cloudFiles") 
  .option("cloudFiles.format", "csv")
  .option("header", "true") 
  .option("cloudFiles.includeExistingFiles", "true") 
  .option("cloudFiles.inferColumnTypes", "true") 
  .option("cloudFiles.schemaLocation", checkpoint_path)
  # 正确指定Autoloader文件名过滤规则
  .option("cloudFiles.pathGlobFilter", "String_Pattern_*.csv")
  # 明确开启递归扫描子目录(可选,默认已启用)
  .option("recursiveFileLookup", "true")
  .load("abfss://<container-path>/Dir1/Dir2/")
)

方案2:通过路径通配符匹配目标文件

直接在加载路径中使用/**/通配符匹配所有子目录下的目标文件,无需额外配置过滤参数:

df = (spark.readStream
  .format("cloudFiles") 
  .option("cloudFiles.format", "csv")
  .option("header", "true") 
  .option("cloudFiles.includeExistingFiles", "true") 
  .option("cloudFiles.inferColumnTypes", "true") 
  .option("cloudFiles.schemaLocation", checkpoint_path)
  # 路径中使用通配符匹配所有子目录下的目标文件
  .load("abfss://<container-path>/Dir1/Dir2/**/String_Pattern_*.csv")
)

补充:手动指定Schema避免推断失败

若仍出现Schema推断错误,可手动定义Schema并传入:

from pyspark.sql.types import StructType, StructField, StringType, IntegerType

# 根据实际CSV结构定义Schema
custom_schema = StructType([
    StructField("column1", StringType(), nullable=True),
    StructField("column2", IntegerType(), nullable=True),
    # 添加更多字段...
])

df = (spark.readStream
  .format("cloudFiles") 
  .option("cloudFiles.format", "csv")
  .option("header", "true") 
  .option("cloudFiles.includeExistingFiles", "true") 
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .option("cloudFiles.pathGlobFilter", "String_Pattern_*.csv")
  .schema(custom_schema)
  .load("abfss://<container-path>/Dir1/Dir2/")
)

注意事项

  • 检查文件名大小写是否与过滤规则一致,Azure Blob存储对文件名大小写敏感;
  • 若之前的checkpoint目录存在旧状态,建议清理后重新运行,避免旧状态干扰新的过滤逻辑;
  • 确保Databricks Runtime版本支持相关Autoloader参数(建议使用DBR 10.0及以上版本)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 02:48:16