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

如何在ADF PySpark Notebook中检查ADLS Gen2目录是否存在

解决PySpark中ADLS Gen2目录存在性判断及数据合并问题

问题场景

在PySpark Notebook处理数据时,需实现以下逻辑:

  • 检查Azure Data Lake Storage Gen2容器中指定目录是否存在
  • 目录存在时:读取目录下Parquet文件为dfStaging
    • 若dfStaging无数据,直接将dfLanding写入该目录
    • 若dfStaging有数据,将dfLanding与dfStaging合并后写入
  • 目录不存在时:直接将dfLanding写入目标目录

原代码在目录不存在时执行spark.read.load会报错,核心问题是缺少目录存在性前置判断。

解决方案

利用Hadoop文件系统API实现ADLS Gen2路径的存在性检查,再分分支处理数据逻辑,同时修正原代码中dfPartition的笔误(应为dfLanding)。

修改后的完整代码

from pyspark.sql import functions as F
from pyspark.sql.window import Window
from org.apache.hadoop.fs import Path

for row in collectPartition: 
    # 筛选对应分区的Landing数据
    dfLanding = df.filter(F.col("_partition") == row["_partition"])
    
    # 获取Hadoop文件系统实例,用于检查目录存在性
    hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration()
    target_path = Path(row["_location"])
    fs = target_path.getFileSystem(hadoop_conf)
    
    # 目录存在的情况
    if fs.exists(target_path):
        dfStaging = spark.read.load(row["_location"], format='parquet', schema=dataSchema)
        
        # 判断Staging数据是否为空
        if dfStaging.count() == 0:
            # 首次写入:覆盖写入Landing数据
            dfLanding.write.mode('overwrite').parquet(row["_location"])
        else:
            # 合并数据:标记数据源,保留最新RECID数据
            dfLanding = dfLanding.withColumn('_source', F.lit(1))  # 标记为新数据
            dfStaging = dfStaging.withColumn('_source', F.lit(2))  # 标记为旧数据
            
            # 合并两个数据集(修正原代码笔误:dfPartition改为dfLanding)
            dfMerge = dfLanding.union(dfStaging)
            
            # 按RECID分区,按_source排序(新数据优先),生成版本号
            partition_window = Window.partitionBy("RECID").orderBy(F.col("_source"))    
            dfMerge = dfMerge.withColumn("_version", F.row_number().over(partition_window))
            
            # 保留最新版本数据,清理辅助字段
            dfMerge = dfMerge.filter(F.col("_version") == 1)
            dfMerge = dfMerge.drop("_source", "_version")
            
            # 覆盖写入合并后的数据
            dfMerge.write.mode('overwrite').parquet(row["_location"])
    else:
        # 目录不存在:直接写入Landing数据
        dfLanding.write.mode('overwrite').parquet(row["_location"])

关键说明

  1. 目录存在性判断:通过Hadoop的FileSystem API直接检查ADLS Gen2路径,避免了直接读取不存在路径引发的异常
  2. 数据合并逻辑:通过_source标记新旧数据,结合窗口函数row_number确保相同RECID保留最新的Landing数据
  3. 笔误修正:原代码中dfPartition.union(dfStaging)应为dfLanding.union(dfStaging),否则会引发变量未定义错误

内容的提问来源于stack exchange,提问作者LordRofticus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 10:03:13