如何在Databricks的PySpark中检查文件/文件夹存在性以避免异常
解决Databricks PySpark中文件存在性检查的问题
你遇到的问题很典型——PySpark的load()方法抛出的异常并不是IOError,而是pyspark.sql.utils.AnalysisException,所以你的except IOError分支根本不会被触发。下面给你几种在Databricks环境下可靠的解决方案:
方法1:捕获正确的异常类型
直接捕获AnalysisException,并通过异常信息判断是否是文件不存在的情况:
from pyspark.sql import SparkSession from pyspark.conf import SparkConf from pyspark.sql.utils import AnalysisException # 初始化SparkSession(注意要调用getOrCreate()获取实例) spark = SparkSession.builder.config(conf=SparkConf()).getOrCreate() try: df = spark.read.format('com.databricks.spark.csv')\ .option("delimiter", ",")\ .options(header='true', inferschema='true')\ .load('/FileStore/tables/HealthCareSample_dumm.csv') print("File Exists") except AnalysisException as e: # 判断异常信息中是否包含路径不存在的关键词 if "Path does not exist" in str(e): print("File not found") else: # 如果是其他类型的分析异常,重新抛出不吞掉 raise e
方法2:使用Databricks的dbutils主动检查(推荐)
Databricks提供了dbutils.fs工具,可以直接检查文件/目录是否存在,这种方法更主动,避免触发异常:
from pyspark.sql import SparkSession from pyspark.conf import SparkConf spark = SparkSession.builder.config(conf=SparkConf()).getOrCreate() file_path = '/FileStore/tables/HealthCareSample_dumm.csv' # 检查文件是否存在:dbutils.fs.ls会返回路径下的文件列表,不存在则抛出异常,所以用try-except包裹 def check_file_exists(file_path): try: # 如果路径存在,ls会返回非空列表 files = dbutils.fs.ls(file_path) return len(files) > 0 except Exception: return False if check_file_exists(file_path): df = spark.read.format('com.databricks.spark.csv')\ .option("delimiter", ",")\ .options(header='true', inferschema='true')\ .load(file_path) print("File Exists") else: print("File not found")
方法3:使用Hadoop FileSystem API(通用跨环境)
如果需要兼容非Databricks的Spark环境,可以直接调用Hadoop的FileSystem API来检查文件:
from pyspark.sql import SparkSession from pyspark.conf import SparkConf from py4j.java_gateway import java_import spark = SparkSession.builder.config(conf=SparkConf()).getOrCreate() # 导入Hadoop的FileSystem和Path类 java_import(spark._jvm, 'org.apache.hadoop.fs.FileSystem') java_import(spark._jvm, 'org.apache.hadoop.fs.Path') file_path = '/FileStore/tables/HealthCareSample_dumm.csv' # 获取Hadoop文件系统实例 fs = spark._jvm.FileSystem.get(spark._jsc.hadoopConfiguration()) if fs.exists(spark._jvm.Path(file_path)): df = spark.read.format('com.databricks.spark.csv')\ .option("delimiter", ",")\ .options(header='true', inferschema='true')\ .load(file_path) print("File Exists") else: print("File not found")
方法对比
- 方法1:代码改动最小,但依赖异常信息的字符串内容,若Spark版本更新异常信息变化可能失效。
- 方法2:Databricks环境下最推荐,简单直观,利用平台原生工具。
- 方法3:通用性最强,适用于任何Spark环境,但代码稍复杂。
内容的提问来源于stack exchange,提问作者Amareshwar Reddy
相关产品推荐
相关产品推荐

