Azure Synapse中PySpark基于DataFrame匹配ADLS Gen2指定文件并处理
解决方案
1. 提取目标国家标识
从现有DataFrame df 中取出唯一的COUNTRY_NAME值(因列内为统一值,无需处理多值情况):
country_id = df.select("COUNTRY_NAME").distinct().first()[0]
2. 定位ADLS文件路径并获取文件列表
先构建ADLS Gen2的目标文件夹路径,再拉取该路径下所有CSV文件的完整路径:
# 替换为你的存储账户名称 storage_account = "your-storage-account" base_path = f"abfss://DETAILS@{storage_account}.dfs.core.windows.net/COUNTRIES DETAIL/YEAR/" # Databricks环境下用dbutils获取文件列表 file_paths = [f.path for f in dbutils.fs.ls(base_path) if f.path.endswith(".csv")] # 非Databricks环境用Hadoop FS API替代 # from pyspark.sql import SparkSession # spark = SparkSession.builder.getOrCreate() # fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration()) # hadoop_path = spark._jvm.org.apache.hadoop.fs.Path(base_path) # file_paths = [str(f.getPath()) for f in fs.listFiles(hadoop_path, False) if str(f.getPath()).endswith(".csv")]
3. 筛选符合匹配规则的文件
根据文件名与COUNTRY_NAME的对应关系筛选文件:
filtered_files = [] for path in file_paths: # 剥离路径和后缀,获取纯文件名 filename = path.split("/")[-1].removesuffix(".csv") # 提取文件名中DETAILS之后的国家相关部分(比如YYYY_DETAILS_INDIA_GOOD_ → INDIA_GOOD) file_country_segment = filename.split("DETAILS_")[-1].rstrip("_") # 去掉下划线后和COUNTRY_NAME对比,匹配则保留路径 if file_country_segment.replace("_", "") == country_id: filtered_files.append(path)
4. 读取匹配文件到DataFrame
如果找到匹配文件,直接读取到df2;否则输出提示:
if filtered_files: df2 = spark.read.csv(filtered_files, header=True, inferSchema=True) # 这里写后续转换处理逻辑 else: print(f"No files matched for country identifier: {country_id}")
关键说明
- 先筛选路径再读取,比全量读取后过滤更高效,减少不必要的IO开销。
- 匹配逻辑可根据实际文件名规则调整:比如如果
COUNTRY_NAME是文件名中某个固定位置的子串,可直接用in关键字简化判断。
内容的提问来源于stack exchange,提问作者BigData Lover
相关产品推荐
相关产品推荐

