Databricks写入CSV报错:无法推断str类型的Schema
问题:Databricks写入自定义名称CSV并使用CRLF时出现schema推断错误
我在Databricks中尝试将DataFrame写入指定名称的CSV文件(替换默认的part-****命名),并要求使用CRLF换行符。为此编写了如下函数:
def WriteOutputFiles (DF_To_Write, FileName, FileVers): TempDir = "/PathToTempDir/" FinalOutputPath = '/PathToFinalDir/' + FileVers + '/' + FileVers + " OutputFiles/" + FileName DF_To_Write.coalesce(1).write.mode("overwrite").option("header", "true").csv(TempDir) #List the files in the temporary directory to find the CSV file first time CSV_Files_First = [file.path for file in dbutils.fs.ls(TempDir) if file.name.endswith(".csv")] if CSV_Files_First: CSV_File = CSV_Files_First[0] #Reading csv file to add CR LF line terminators Temp_DF = spark.read.option("header", "true").csv(CSV_File) Final_DF_CRLF = Temp_DF.rdd.map(lambda row: ','.join([str(elem) for elem in row]) + '\r\n').toDF(['All Data']) Final_DF_CRLF.coalesce(1).write.mode("overwrite").option("header", "true").csv(TempDir) #List the files in the temporary directory to find the CSV file second time CSV_Files_Second = [file.path for file in dbutils.fs.ls(TempDir) if file.name.endswith(".csv")] if CSV_Files_Second: CSV_File = CSV_Files_Second[0] #Move and rename the CSV file to the desired location dbutils.fs.mv(CSV_File, FinalOutputPath) #Remove the temporary directory dbutils.fs.rm(TempDir, true)
调用函数时出现错误Can not infer schema for type: 'str',问题出在这行代码:
Final_DF_CRLF = Temp_DF.rdd.map(lambda row: ','.join([str(elem) for elem in row]) + '\r\n').toDF(['All Data'])
解决方案
错误原因分析
当你将RDD映射为单个字符串后,调用toDF(['All Data'])时,Spark无法自动推断该字符串的结构化schema。toDF()需要明确知道列的数据类型,而纯字符串类型无法让Spark生成有效的列定义,因此抛出推断失败的错误。
方案1:手动指定schema修复现有代码
通过手动定义单列schema,使用createDataFrame替代toDF,明确告诉Spark列的类型:
from pyspark.sql.types import StringType, StructField, StructType def WriteOutputFiles (DF_To_Write, FileName, FileVers): TempDir = "/PathToTempDir/" FinalOutputPath = '/PathToFinalDir/' + FileVers + '/' + FileVers + " OutputFiles/" + FileName DF_To_Write.coalesce(1).write.mode("overwrite").option("header", "true").csv(TempDir) #List the files in the temporary directory to find the CSV file first time CSV_Files_First = [file.path for file in dbutils.fs.ls(TempDir) if file.name.endswith(".csv")] if CSV_Files_First: CSV_File = CSV_Files_First[0] #Reading csv file to add CR LF line terminators Temp_DF = spark.read.option("header", "true").csv(CSV_File) # 手动定义单列schema single_col_schema = StructType([StructField("All Data", StringType(), nullable=True)]) # 使用createDataFrame替代toDF,传入定义好的schema Final_DF_CRLF = spark.createDataFrame( Temp_DF.rdd.map(lambda row: ','.join([str(elem) for elem in row]) + '\r\n'), schema=single_col_schema ) Final_DF_CRLF.coalesce(1).write.mode("overwrite").option("header", "true").csv(TempDir) #List the files in the temporary directory to find the CSV file second time CSV_Files_Second = [file.path for file in dbutils.fs.ls(TempDir) if file.name.endswith(".csv")] if CSV_Files_Second: CSV_File = CSV_Files_Second[0] #Move and rename the CSV file to the desired location dbutils.fs.mv(CSV_File, FinalOutputPath) #Remove the temporary directory dbutils.fs.rm(TempDir, True)
方案2:更高效的原生配置方式(推荐)
Spark的CSV Writer本身支持指定换行符,无需通过RDD转换手动添加CRLF,也不需要两次读写临时文件,大幅提升效率:
def WriteOutputFiles(DF_To_Write, FileName, FileVers): TempDir = "/PathToTempDir/" # 使用f-string简化路径拼接 FinalOutputPath = f'/PathToFinalDir/{FileVers}/{FileVers} OutputFiles/{FileName}' # 直接写入临时目录,通过lineSep参数指定CRLF换行符 DF_To_Write.coalesce(1).write.mode("overwrite")\ .option("header", "true")\ .option("lineSep", "\r\n")\ # 核心配置:指定CRLF作为换行符 .csv(TempDir) # 找到生成的CSV文件 CSV_Files = [file.path for file in dbutils.fs.ls(TempDir) if file.name.endswith(".csv")] if CSV_Files: CSV_File = CSV_Files[0] # 移动并重命名到目标路径 dbutils.fs.mv(CSV_File, FinalOutputPath) # 删除临时目录(递归删除所有内容) dbutils.fs.rm(TempDir, True)
这个方案避免了重复IO和RDD转换的开销,同时完全满足自定义文件名和CRLF换行的需求。
内容的提问来源于stack exchange,提问作者Danylo Kuznetsov
相关产品推荐
相关产品推荐

