Databricks生成Parquet文件报错:路径不存在(实际文件存在)
问题描述
在Databricks中尝试将CSV文件转换为Parquet文件,已确认输入目录及目标CSV文件存在且路径正确,但始终报错提示第一个CSV文件路径不存在,无法继续操作。相关代码如下:
import os from pyspark.sql.types import StructType, StructField, StringType # Define the schema for the files you want to convert from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType, DoubleType schema = StructType([ StructField("METER_ADDRESS", StringType(), True), StructField("READING_DATE", TimestampType(), True), StructField("READING_VALUE_L", DoubleType(), True), StructField("LOW_BATTERY_ALR", IntegerType(), True), StructField("LEAK_ALR", IntegerType(), True), StructField("MAGNETIC_TAMPER_ALR", IntegerType(), True), StructField("METER_ERROR_ALR", IntegerType(), True), StructField("BACK_FLOW_ALR", IntegerType(), True), StructField("BROKEN_PIPE_ALR", IntegerType(), True), StructField("EMPTY_PIPE_ALR", IntegerType(), True), StructField("SPECIFIC_ERROR_ALR", IntegerType(), True) ]) # Set the input and output directories input_directory = "/dbfs/FileStore/tables/Calybre Capstone Project - Part 1" output_directory = "dbfs:/FileStore/tables/Calybre Capstone Project - Part 1/Parquet Files" # Iterate over each file in the input directory for filename in os.listdir(input_directory): if filename.endswith(".csv"): filepath = os.path.join(input_directory, filename) # Read in the file using spark.read() df = spark.read.csv(filepath, header=True, schema=schema) # Write the resulting DataFrame as a parquet file output_path = os.path.join(output_directory, filename + ".parquet") df.write.parquet(output_path)
问题原因
- 本地与分布式文件操作混用:
os.listdir是Driver节点本地文件系统操作,仅能识别Driver节点本地挂载的DBFS内容,若目录存在子目录,会被误判为CSV文件,拼接路径后Spark读取时报错;同时该操作无法适配Databricks分布式存储的特性,可能出现路径识别偏差。 - 路径格式不统一:输入目录用
/dbfs/格式,输出目录用dbfs:/格式,虽然Databricks兼容两种格式,但混用可能引发路径解析冲突。
解决方案
改用Databricks原生的dbutils.fs工具遍历文件,统一路径格式,避免本地操作与分布式操作的冲突:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType, DoubleType # 定义Schema schema = StructType([ StructField("METER_ADDRESS", StringType(), True), StructField("READING_DATE", TimestampType(), True), StructField("READING_VALUE_L", DoubleType(), True), StructField("LOW_BATTERY_ALR", IntegerType(), True), StructField("LEAK_ALR", IntegerType(), True), StructField("MAGNETIC_TAMPER_ALR", IntegerType(), True), StructField("METER_ERROR_ALR", IntegerType(), True), StructField("BACK_FLOW_ALR", IntegerType(), True), StructField("BROKEN_PIPE_ALR", IntegerType(), True), StructField("EMPTY_PIPE_ALR", IntegerType(), True), StructField("SPECIFIC_ERROR_ALR", IntegerType(), True) ]) # 统一使用dbfs:/格式路径 input_directory = "dbfs:/FileStore/tables/Calybre Capstone Project - Part 1" output_directory = "dbfs:/FileStore/tables/Calybre Capstone Project - Part 1/Parquet Files" # 遍历目录下的CSV文件,排除子目录 for file_info in dbutils.fs.ls(input_directory): if not file_info.isDir() and file_info.name.endswith(".csv"): filepath = file_info.path # 读取CSV文件 df = spark.read.csv(filepath, header=True, schema=schema) # 构造输出路径 output_filename = file_info.name.replace(".csv", ".parquet") output_path = f"{output_directory}/{output_filename}" # 写入Parquet文件,存在则覆盖 df.write.mode("overwrite").parquet(output_path)
额外说明
dbutils.fs.ls可正确识别DBFS分布式文件,返回的file_info包含路径、是否为目录等信息,能有效避免误处理子目录。mode("overwrite")用于避免重复运行时因输出路径已存在报错,可根据需求改为mode("append")或mode("ignore")。
内容的提问来源于stack exchange,提问作者Muhammed Rif'at Kader
相关产品推荐
相关产品推荐

