如何使用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/") )
原因分析
- 参数名错误:使用了Spark原生的
pathGlobfilter(小写f)而非Autoloader专用的cloudFiles.pathGlobFilter参数,导致过滤规则未被正确识别; - 过滤规则误用:
pathGlobFilter仅用于匹配文件名,不支持路径通配符(如/**/),错误的规则会导致没有匹配到任何文件; - 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
相关产品推荐
相关产品推荐

