You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark读取含多嵌套分隔符的管道符分隔CSV文件失败问题咨询

Spark嵌套分隔符CSV解析解决方案

问题根因

  • 首次使用分号作为CSV分隔符读取失败:文件表头为|分隔,与数据行按分号拆分的列数不匹配,Spark CSV解析时的结构校验不通过触发IO异常
  • 后续使用|作为分隔符读取结果不符合预期:u5字段的pid值本身包含|分隔符,会被误拆分为多列,破坏原有数据结构

实现逻辑

不要直接使用CSV API按分隔符拆分,先整行读取文本后自定义嵌套结构解析逻辑:

  1. 整行读取文件,过滤表头行
  2. 每行仅按第一个|拆分,得到独立的Ord_value字段和后续所有配置字段串
  3. 将配置字段串按;拆分提取U变量键值对,存入Map
  4. 分别提取单值字段、拆分列表类字段:pid按|拆分,其余列表字段按,拆分
  5. 用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 06:18:03