Spark解析无引号包裹的CSV时如何移除列内换行符生成正确DataFrame
无引号包裹字段换行的CSV解析方案
问题场景
解析CSV文件生成DataFrame时,若部分列存在未被引号包裹的换行符(\n或\r),常规.option("multiline",true)参数无法生效,会出现字段错位、空值的问题。
问题示例
原始CSV内容
id,commment,name 1,good,bob 2,bad ,tim 3,fine,sarah
直接读取的错误输出
| id | comment | name |
|---|---|---|
| 1 | good | bob |
| 2 | bad | null |
| tim | null | null |
| 3 | fine | sarah |
期望输出
| id | comment | name |
|---|---|---|
| 1 | good | bob |
| 2 | bad | tim |
| 3 | fine | sarah |
实现方案
由于未加引号的字段换行无法被Spark CSV的内置规则识别为字段内部换行,需要先做文本预处理,再转结构化数据,具体实现如下:
实现逻辑
- 先以文本格式读取原始CSV,保留所有原始行内容
- 按表头确定CSV总列数,逐行拼接内容,直到拼接后的内容按分隔符拆分的字段数等于总列数,即为一条完整记录
- 将处理后的完整记录拆分字段,得到目标DataFrame
Scala示例代码
// 读取原始文本RDD val rawRDD = spark.sparkContext.textFile("test.csv") // 获取表头和总列数 val header = rawRDD.first() val colTotal = header.split(",").length // 预处理拼接换行的不完整行 val fixedRDD = rawRDD .filter(_ != header) .mapPartitions(iter => { var cache = "" iter.flatMap(currentLine => { cache += currentLine val fields = cache.split(",", -1) if (fields.length >= colTotal) { val fullLine = cache cache = "" List(fullLine) } else { Nil } }) }) // 转DataFrame并拆分字段 import spark.implicits._ val resultDF = fixedRDD.toDF("line") .selectExpr( "split(line, ',')[0] as id", "split(line, ',')[1] as comment", "split(line, ',')[2] as name" ) // 输出验证 resultDF.show()
执行输出
+---+-------+-----+ | id|comment| name| +---+-------+-----+ | 1| good| bob| | 2| bad| tim| | 3| fine|sarah| +---+-------+-----+
内容的提问来源于stack exchange,提问作者granpawalton
相关产品推荐
相关产品推荐

