Fabric中PySpark如何列出数据湖目录并获取最新目录
解决方案:在Fabric PySpark中找到最新的Data子目录
针对你在Fabric环境下的需求,以下两种方案可以实现列出Files/Landing下的DataYYYYMMDDHHMM格式子目录并筛选最新项的功能:
方案一:使用mssparkutils直接操作文件系统(推荐)
Fabric内置的mssparkutils工具可以直接与数据湖交互,避免Hadoop FS路径解析的问题,效率更高:
# 替换为你的基础目录路径 base_path = "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/Files/Landing/" # 列出基础目录下的所有子项 dir_items = mssparkutils.fs.ls(base_path) # 筛选出以Data开头的子目录 data_dirs = [item for item in dir_items if item.isDir and item.name.startswith("Data")] # 按目录名称逆序排序(因命名格式为DataYYYYMMDDHHMM,字符串排序即可识别最新) data_dirs_sorted = sorted(data_dirs, key=lambda x: x.name, reverse=True) # 获取最新的目录对象及完整路径 latest_dir = data_dirs_sorted[0] latest_dir_path = latest_dir.path # 后续可读取该目录下的CSV文件 df_latest = spark.read.csv(latest_dir_path, header=True, inferSchema=True)
方案二:纯Spark API实现(无需依赖mssparkutils)
如果不想使用工具类,可以通过Spark的文件元数据提取目录信息:
from pyspark.sql.functions import input_file_name, split, element_at base_path = "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/Files/Landing/" # 读取所有子目录下的CSV文件(仅需获取路径,可按需调整读取方式) df = spark.read.csv(base_path + "*/*", header=True, inferSchema=True) # 添加文件路径列,并提取父目录名(即Data开头的子目录) df_dir_info = df.withColumn("file_path", input_file_name()) \ .withColumn("dir_name", element_at(split("file_path", "/"), -2)) # 获取所有唯一的Data子目录名称 unique_dirs = df_dir_info.select("dir_name").distinct().rdd.flatMap(lambda x: x).collect() # 排序后取最新目录 latest_dir = sorted(unique_dirs, reverse=True)[0] latest_dir_path = base_path + latest_dir + "/" # 读取最新目录下的文件进行转换 df_latest = spark.read.csv(latest_dir_path, header=True, inferSchema=True)
关于之前报错的说明
你使用Hadoop FS API时出现的FriendlyNameSupportDisabled错误,是因为Fabric数据湖的路径解析逻辑与原生Hadoop FS存在差异,自动添加的GUID导致路径无法被正确识别。使用上述两种方案可以规避该问题。
内容的提问来源于stack exchange,提问作者Gav Cheal
相关产品推荐
相关产品推荐

