使用Spark DataFrame与Scala无法生成CSV文件,反而创建同名文件夹
问题分析与解决方案
核心原因
Spark作为分布式计算框架,默认会将输出数据按分区拆分为多个文件,而非单个文件。当你指定类似"C:\Users\name\output.csv"的路径时,Spark会将该路径识别为文件夹名称,而非文件名,因此会创建对应文件夹并尝试在其中写入分区文件。另外,你的路径使用了单反斜杠,在Scala中属于转义字符,可能导致路径解析异常,进而出现临时文件夹无内容的情况。
解决步骤
1. 修正路径转义问题
Scala中路径的反斜杠需要转义,或者使用原始字符串:
// 方式1:双反斜杠转义 val outputDir = "C:\\Users\\name\\temp_output" // 方式2:原始字符串(Scala 2.13+支持) val outputDir = raw"C:\Users\name\temp_output"
2. 输出单个CSV文件的两种方式
方式一:合并分区后输出(推荐)
使用coalesce(1)将数据合并到一个分区(避免shuffle,性能更优),输出到临时文件夹后,手动将分区文件重命名为目标CSV:
import java.io.File import org.apache.commons.io.FileUtils // 定义临时输出文件夹和目标文件名 val outputDir = "C:\\Users\\name\\temp_output" val targetFileName = s"C:\\Users\\name\\file_${timestamp}.csv" // 写入临时文件夹 dataResults .withColumn("primaryid", concat_ws("$", col("primaryid"))) .coalesce(1) // 合并为1个分区 .write .option("delimiter", "\t") .option("header", "true") .mode(SaveMode.Overwrite) .csv(outputDir) // Spark 2.0+内置CSV支持,无需依赖com.databricks.spark.csv // 找到分区文件并重命名 val partFile = new File(outputDir).listFiles() .filter(_.getName.startsWith("part-")) .headOption.getOrElse(throw new Exception("未找到输出的分区文件")) FileUtils.moveFile(partFile, new File(targetFileName)) // 删除临时文件夹 FileUtils.deleteDirectory(new File(outputDir))
方式二:自定义表头与RDD写入(适用于简单格式)
如果数据格式简单,可将DataFrame转为RDD后自行处理表头:
// 生成表头和数据RDD val header = "primaryid" val dataRDD = dataResults .withColumn("primaryid", concat_ws("$", col("primaryid"))) .rdd.map(_.mkString("\t")) // 合并表头与数据 val finalRDD = sqlContext.sparkContext.parallelize(Seq(header)) ++ dataRDD finalRDD.coalesce(1).saveAsTextFile(outputDir) // 后续重命名、删除临时文件夹步骤同方式一
3. 注意事项
- 生产环境处理大规模数据时,避免使用
coalesce(1)或repartition(1),会将所有数据集中到单个节点,可能引发内存溢出,仅适用于小数据量场景。 - 确认
SaveMode.Overwrite生效,避免输出路径被其他进程占用导致写入失败。
内容的提问来源于stack exchange,提问作者Braham Shakti
相关产品推荐
相关产品推荐

