Spark读取含多嵌套分隔符的管道符分隔CSV文件失败问题咨询
Spark嵌套分隔符CSV解析解决方案
问题根因
- 首次使用分号作为CSV分隔符读取失败:文件表头为|分隔,与数据行按分号拆分的列数不匹配,Spark CSV解析时的结构校验不通过触发IO异常
- 后续使用|作为分隔符读取结果不符合预期:u5字段的pid值本身包含|分隔符,会被误拆分为多列,破坏原有数据结构
实现逻辑
不要直接使用CSV API按分隔符拆分,先整行读取文本后自定义嵌套结构解析逻辑:
- 整行读取文件,过滤表头行
- 每行仅按第一个|拆分,得到独立的Ord_value字段和后续所有配置字段串
- 将配置字段串按;拆分提取U变量键值对,存入Map
- 分别提取单值字段、拆分列表类字段:pid按|拆分,其余列表字段按,拆分
- 用
arrays_zip将所有列表字段按位置对齐,再通过explode函数将数组拆分为多行
完整可运行Scala代码
import org.apache.spark.sql.SparkSession import org.apache.spark.SparkConf import org.apache.spark.sql.functions._ object example3 { def main(args:Array[String]):Unit={ System.setProperty("hadoop.home.dir", "C:\\hadoop\\") val conf = new SparkConf().setAppName("parse_demo").setMaster("local[*]") val spark = SparkSession.builder().config(conf).getOrCreate() spark.sparkContext.setLogLevel("Error") import spark.implicits._ // 1. 整行读取文本,过滤表头 val rawData = spark.read.textFile("file:///C:/Users/User/Desktop/example3.txt") .filter(!_.startsWith("Ord_value")) // 2. 解析每行数据提取字段 val parsedDF = rawData.map(line => { // 只按第一个|拆分,避免后续pid中的|被误拆 val Array(ordValue, otherPart) = line.split("\\|", 2) // 拆分得到所有U变量键值对 val kvMap = otherPart.split(";") .filter(_.nonEmpty) .map(kvStr => { val Array(k, v) = kvStr.split("=", 2) k.toLowerCase -> v }).toMap // 提取各字段,列表类按对应分隔符拆分 val orderId = kvMap("u1") val pidArr = kvMap("u5").split("\\|") val nameArr = kvMap("u15").split(",") val priceArr = kvMap("u16").split(",") val quantityArr = kvMap("u17").split(",") (ordValue, orderId, pidArr, nameArr, priceArr, quantityArr) }).toDF("Ord_value", "orderid", "pid_arr", "name_arr", "price_arr", "quantity_arr") // 3. 数组对齐后炸开为多行 val resultDF = parsedDF.select( col("Ord_value"), col("orderid"), explode(arrays_zip(col("pid_arr"), col("name_arr"), col("price_arr"), col("quantity_arr"))).as("item") ).select( col("Ord_value"), col("orderid"), col("item.pid_arr").as("pid"), col("item.name_arr").as("name"), col("item.price_arr").as("price"), col("item.quantity_arr").as("quantity") ) // 打印结果 resultDF.show(false) // 如需输出为|分隔的CSV文件,执行以下代码 // resultDF.write.option("delimiter", "|").option("header", "true").csv("你的输出目录路径") spark.stop() } }
注意事项
arrays_zip为Spark 2.4及以上版本内置函数,若使用更低版本Spark,可自定义UDF实现多数组按位置对齐逻辑- 拆分字符串时指定拆分次数为2,可避免字段值中包含分隔符时被误拆
- 若实际数据中存在U变量缺失的情况,可在代码中增加空值判断逻辑避免数组越界
内容的提问来源于stack exchange,提问作者Anonymous
相关产品推荐
相关产品推荐

