Spark读取含换行与双引号转义的CSV时multiLine失效问题
解决Spark读取CSV时multiLine配置失效的问题
问题根源
你当前的实现先把文件按行拆分成RDD/Dataset,每个元素是单独一行文本。此时Spark的CSV读取器会将每个元素视为完整的CSV记录,multiLine选项的作用是读取文件时识别跨多行的记录,而非处理已经拆分好的单行数据集,因此即使开启该选项也无法合并分散的行来补全被双引号包裹的多行值。
最优解决方案:使用Spark CSV原生参数
Spark CSV读取器原生支持处理双引号转义和多行记录,无需手动读取文本替换,直接配置escape参数即可:
def readCSVWithProperOptions(path: String): DataFrame = { spark.read .option("header", true) .option("multiLine", true) .option("escape", "\"") // 用双引号作为转义符,自动解析""为单个" .csv(path) }
该方案会自动完成:
- 将值内的
""解析为单个" - 正确识别被双引号包裹的跨多行字段,不会拆分记录
自定义文本处理场景的解决方案
如果必须先做自定义文本替换,需要先将属于同一条记录的多行合并为单个元素,再交给CSV读取器:
def readTextReplaceQuotesAndParseAsCSV(path: String): DataFrame = { import org.apache.spark.sql.Encoders val rdd = spark.sparkContext.textFile(path) // 合并同一条记录的多行:通过双引号数量判断记录是否结束 val mergedRDD = rdd.mapPartitions { lines => val recordBuffer = new StringBuilder lines.flatMap { line => val quoteCount = line.count(_ == '"') recordBuffer.append(line) // 双引号数量为偶数,说明当前记录完整,输出并清空缓冲区 if (quoteCount % 2 == 0) { val completeRecord = recordBuffer.toString() recordBuffer.clear() Some(completeRecord) } else { // 双引号数量为奇数,记录未结束,追加换行符后继续合并下一行 recordBuffer.append("\n") None } } ++ { // 处理分区内最后一条未闭合的记录 if (recordBuffer.nonEmpty) Some(recordBuffer.toString()) else None } } // 替换连续双引号为单个双引号 val replacedRDD = mergedRDD.map(_.replaceAll("\"\"", "\"")) val ds = spark.createDataset(replacedRDD)(Encoders.STRING) val df = spark.read .option("header", true) .option("multiLine", true) .csv(ds) df }
验证结果
处理你的测试文件后,会得到正确的解析结果:
+---+-----+--------------------+----+----+ | id| name| desc|tags|user| +---+-----+--------------------+----+----+ | 1| sam| quoted two\nlines|tag1| tom| | 2|chuck|one line, two quotes|tag1| tom| | 3| sam|quoted\nthree\nlines|tag1| tom| +---+-----+--------------------+----+----+
内容的提问来源于stack exchange,提问作者Aravind Yarram
相关产品推荐
相关产品推荐

