Pyspark使用SparkFiles读取Parquet文件时报CSV不存在错误的原因
问题原因分析
1. SparkSession元数据缓存冲突
若当前SparkSession生命周期内,你曾使用
spark.read.csv()读取过/a/b/c.parquet路径(哪怕读取操作未成功),Spark会默认缓存该路径对应的数据源格式、Schema等元数据。后续调用spark.read.parquet()读取同一路径时,Spark可能优先调用缓存的元数据,错误将该路径识别为CSV格式,按CSV规则读取文件,触发CSV相关的文件缺失报错。
2. addFile未开启递归导致目录内容不完整
spark.sparkContext.addFile()的recursive参数默认值为False,如果传入路径是目录,只会分发目录本身,不会递归分发目录下的所有子文件。Parquet格式通常以目录形式存储,包含多个数据分片、元数据文件,未开启递归分发时,executor工作目录下的c.parquet目录为空或内容不完整。部分版本Spark读取不完整的Parquet目录时,会出现格式识别异常,错误匹配到CSV数据源的读取逻辑。
3. Parquet目录下存在CSV格式残留文件
如果
c.parquet目录中存在写入操作残留的CSV格式临时文件、元数据文件(比如错误写入的_SUCCESS.csv、临时分片文件),Spark读取目录时会遍历所有文件,旧版本Spark默认不会自动过滤非Parquet格式文件,会尝试解析所有文件,遇到CSV文件时就会触发CSV解析相关报错。
解决方案
- 先清理SparkSession元数据缓存,执行
spark.catalog.clearCache(),也可直接重启SparkSession后重试读取。 - 调用
addFile时开启递归参数,保证目录下所有文件都被分发到executor:parquet_dir = "/a/b/c.parquet" spark.sparkContext.addFile(parquet_dir, recursive=True) - 读取时添加参数忽略非Parquet文件、损坏文件:
parquet_path = SparkFiles.get("c.parquet") spark.read.option("spark.sql.files.ignoreNonParquetFiles", "true")\ .option("spark.sql.files.ignoreCorruptFiles", "true")\ .parquet(f"file://{parquet_path}") - 检查
/a/b/c.parquet目录下的文件列表,删除非Parquet格式的残留文件。
内容的提问来源于stack exchange,提问作者foorx
相关产品推荐
相关产品推荐

