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

Spark Scala读取逗号分隔CSV时JSON格式列被拆分如何处理

解决方案

方案1:调整CSV读取参数(优先推荐)

之前配置未生效的核心原因是参数匹配度不足,需结合CSV实际转义规则配置:

  • 常规CSV中被双引号包裹的字段,内部的双引号默认用双引号自身转义(即""代表单个"),对应配置如下:
import org.apache.spark.sql.SparkSession

val spark = SparkSession.builder()
  .appName("ReadCsvWithJsonColumn")
  .master("local[*]")
  .getOrCreate()

val df = spark.read
  .option("header", "true") // CSV有表头则保留,无表头设为false
  .option("sep", ",")
  .option("quote", "\"") // 指定字段包裹符为双引号
  .option("escape", "\"") // 指定双引号为转义符,适配字段内的双引号转义
  .option("quoteMode", "MINIMAL") // 仅特殊字段用引号包裹,为默认配置,之前设置ALL易引发适配异常
  .csv("你的CSV文件路径")

如果你的CSV中JSON字段内部的双引号用反斜杠转义(即\"代表单个"),将上面escape参数值改为"\\"即可。


方案2:整行读取后手动拆分(兼容非标准格式CSV)

如果CSV本身不符合标准转义规则,可先读取整行再按固定列数拆分,避免JSON列被误切:

import org.apache.spark.sql.functions._

// 整行读取CSV,每行内容存为单个value字段
val rawDf = spark.read.text("你的CSV文件路径")

// 提取表头,limit设为CSV实际总列数,保证JSON列完整
val header = rawDf.head().getString(0).split(",", limit = 实际总列数)
// 过滤表头行
val dataDf = rawDf.filter(row => row.getString(0) != header.mkString(","))

// 按固定列数拆分行内容
val splitDf = dataDf.withColumn("split_cols", split(col("value"), ",", limit = header.length))

// 映射为独立字段
val finalDf = header.zipWithIndex.foldLeft(splitDf){ (df, colInfo) =>
  df.withColumn(colInfo._1, col("split_cols").getItem(colInfo._2))
}.drop("value", "split_cols")

JSON列结构化解析

拿到完整的JSON列后,可通过from_json函数直接解析为结构化字段:

import org.apache.spark.sql.types._

// 定义JSON结构
val jsonSchema = ArrayType(StructType(Seq(
  StructField("code", StringType),
  StructField("name", StringType),
  StructField("type", StringType)
)))

val parsedDf = finalDf.withColumn("parsed_json", from_json(col("你的JSON列名"), jsonSchema))

内容的提问来源于stack exchange,提问作者Ritesh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 08:30:03