Spark 1.6中使用CSV库写入DataFrame时如何指定换行符?
Spark 1.6 指定CSV换行符的解决方案
嘿,刚好对Spark 1.6的CSV处理这块比较熟,来给你详细说说这个问题~
首先明确一点:Spark 1.6版本不管是使用第三方的spark-csv包(当时CSV还没内置到Spark核心,需要单独引入依赖),还是早期的原生数据源,都没有直接配置换行符的参数,默认就是用\n作为行分隔符。不过我们可以通过两种可行的方案来实现自定义换行符的需求:
方案一:数据预处理 + 事后全局替换
这个方案比较简单易操作,适合大多数场景,步骤如下:
预处理数据,替换内部换行符
先把DataFrame中所有字符串列里的\n替换成一个不会出现在数据中的占位符(比如<NEWLINE>),避免后续替换全局换行符时破坏数据内容:- Scala版本:
import org.apache.spark.sql.functions._ val processedDF = originalDF.select( originalDF.columns.map(colName => regexp_replace(col(colName), "\n", "<NEWLINE>").alias(colName)): _* ) - Python版本:
from pyspark.sql.functions import regexp_replace processedDF = originalDF.select([regexp_replace(col(c), "\n", "<NEWLINE>").alias(c) for c in originalDF.columns])
- Scala版本:
写出CSV文件
写出时指定合适的分隔符(建议选数据中不存在的符号,比如|),如果需要表头可以加上header参数:- Scala版本:
processedDF.write .format("com.databricks.spark.csv") .option("header", "true") // 按需开启表头 .option("delimiter", "|") .option("quote", "\"") // 可选:用引号包裹字段,避免分隔符干扰 .save("/path/to/your/output") - Python版本:
processedDF.write \ .format("com.databricks.spark.csv") \ .option("header", "true") \ .option("delimiter", "|") \ .option("quote", "\"") \ .save("/path/to/your/output")
- Scala版本:
全局替换换行符并恢复内部换行
写出完成后,用命令行工具(比如Linux/macOS的sed,Windows的PowerShell)批量替换文件中的\n为你需要的换行符,同时把占位符换回\n:
比如要换成Windows风格的\r\n,Linux/macOS下执行:sed -i 's/\n/\r\n/g; s/<NEWLINE>/\n/g' /path/to/your/output/*.csv
方案二:自定义RDD输出,直接控制换行符
如果需要更灵活的控制,可以把DataFrame转成RDD,自定义输出逻辑来指定换行符,适合对格式要求极高的场景:
以Scala为例,实现自定义OutputFormat来控制换行符:
import org.apache.spark.SparkContext import org.apache.hadoop.fs.Path import org.apache.hadoop.io.{Text, LongWritable} import org.apache.hadoop.mapred.{FileOutputFormat, TextOutputFormat, JobConf, RecordWriter, Reporter, Progressable} import java.io.DataOutputStream // 自定义OutputFormat,指定换行符 class CustomLineEndingOutputFormat extends TextOutputFormat[LongWritable, Text] { override def getRecordWriter(fs: org.apache.hadoop.fs.FileSystem, job: JobConf, name: String, progress: Progressable): RecordWriter[LongWritable, Text] = { val out: DataOutputStream = fs.create(new Path(job.get("mapred.output.dir"), name), progress) // 替换成你需要的换行符,比如"\r\n" val lineEnding = "\r\n".getBytes("UTF-8") new RecordWriter[LongWritable, Text] { override def write(key: LongWritable, value: Text): Unit = { out.write(value.getBytes, 0, value.getLength) out.write(lineEnding) } override def close(reporter: Reporter): Unit = out.close() } } } // 处理DataFrame并写出 val sc: SparkContext = originalDF.sqlContext.sparkContext // 生成表头行 val header = originalDF.columns.mkString(",") // 生成数据行 val dataRows = originalDF.rdd.map(row => row.mkString(",")) // 合并表头和数据 val fullRows = sc.parallelize(Seq(header)) ++ dataRows // 使用自定义OutputFormat写出 fullRows.saveAsHadoopFile( "/path/to/custom/output", classOf[LongWritable], classOf[Text], classOf[CustomLineEndingOutputFormat] )
注意事项
- 如果数据中包含逗号或其他特殊字符,一定要在写出CSV时开启
quote参数,用引号包裹字段,避免CSV格式错乱。 - 方案一中的占位符要确保不会出现在你的业务数据里,否则会导致替换错误。
内容的提问来源于stack exchange,提问作者Defcon
相关产品推荐
相关产品推荐

