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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 18:03:17