PySpark如何按文件掩码匹配ADLS Gen2文件加载到对应DataFrame
PySpark 按文件名匹配加载ADLS Gen2分区Parquet实现方案
不需要手动写循环逐层遍历年/月/日/时间分区,Spark原生内置的文件匹配参数就能搞定,比用dbutils.fs.ls单线程扫目录性能高一个量级,尤其适配你这种分区层级多、单目录文件量大的场景。
前置准备
先确认Databricks已经拿到ADLS Gen2的访问权限,不管是服务主体OAuth授权、存储账户密钥配置还是挂载外部存储都可以,代码里直接替换成你自己的访问路径即可。
核心实现代码
# 替换成你自己的ADLS Gen2 ARCHIVE目录根路径 base_path = "abfss://<你的容器名>@<你的存储账号名>.dfs.core.chinacloudapi.cn/ARCHIVE/" # 配置需要匹配的文件名关键词和对应DataFrame名,后续新增业务文件直接往字典里加就行 target_tables = { "account": "accountdf", "customer": "customerdf", "sales": "salesdf" } # 通用读取配置 read_options = { "recursiveFileLookup": "true", # 自动递归遍历所有子分区目录,跳过手动逐层遍历的逻辑 "pathGlobFilter": None # 后续动态替换成对应文件名匹配规则,匹配文件名含指定关键词的parquet文件 } # 批量加载所有匹配文件 loaded_dfs = {} for file_keyword, df_name in target_tables.items(): current_opt = read_options.copy() # 文件名掩码规则:匹配任意前缀、含指定关键词、后缀为.parquet的文件,自动跳过_SUCCESS、.crc等临时隐藏文件 current_opt["pathGlobFilter"] = f"*{file_keyword}*.parquet" df = spark.read.format("parquet").options(**current_opt).load(base_path) loaded_dfs[df_name] = df # 直接注册为全局变量,后续可以直接调用accountdf/customerdf等变量 locals()[df_name] = df print(f"{file_keyword}文件加载完成,共{df.count()}行,{len(df.columns)}个字段")
常见场景优化提示
- 如果你的分区目录是Hive标准格式(比如
YEAR=2024/Month=05/Day=20/Time=143000),把recursiveFileLookup配置替换为"mergeSchema": "true",Spark会自动把分区字段解析为DataFrame的列,后续和Synapse做数据比对时可以直接用分区字段做时间过滤,不用额外解析路径。 - 如果分区目录是纯数字命名的非标准格式(比如
2024/05/20/143000),需要分区字段的话可以加一列取文件路径自行拆分:from pyspark.sql.functions import input_file_name, split accountdf = accountdf.withColumn("file_path", input_file_name()) # 按路径结构拆分提取年/月/日/时间即可 - 如果不需要加载全量历史分区,比如只需要加载指定时间范围的数据,不用扫整个ARCHIVE根目录,直接在
load()方法里传入限定的路径前缀列表即可,pathGlobFilter会在限定路径范围内匹配文件,扫描速度快很多:# 示例:只加载2024年5月的所有account文件 accountdf_may2024 = spark.read.parquet( base_path + "2024/05/*/*/", pathGlobFilter="*account*.parquet" ) - 后续和Synapse池做数据比对时,建议先按主键对加载的DataFrame做去重,再读取Synapse侧对应时间分区的数据做join校验,避免全表扫描Synapse影响性能。
内容的提问来源于stack exchange,提问作者user774952
相关产品推荐
相关产品推荐

