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

Spark读取含格式不规范JSON的CSV文件问题求助

解决Spark读取含内嵌逗号的JSON列的CSV文件问题

你的CSV文件最后一列是包含逗号的JSON格式字符串,Spark默认按逗号分割会把JSON里的逗号当成列分隔符,导致截断。下面提供两种可靠的解决方法:

方法一:正则拆分非大括号内的逗号

利用正则表达式精准匹配不在大括号内部的逗号作为列分隔符,直接拆分每行数据:

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

// 1. 按文本文件读取原始数据,每行作为单个字符串列
val rawDF = spark.read.text("your_file_path.csv")

// 2. 用正则拆分出4个字段(匹配不在大括号内的逗号)
val splitDF = rawDF.withColumn(
    "parts",
    split(col("value"), ",(?=(?:[^{}]*{[^{}]*})*[^{}]*$)")
).select(
    col("parts")(0).cast(IntegerType).as("c1"),
    col("parts")(1).cast(IntegerType).as("c2"),
    col("parts")(2).cast(IntegerType).as("c3"),
    col("parts")(3).as("json_col")
)

// 3. 定义JSON结构的Schema,解析JSON列
val jsonSchema = StructType(Seq(
    StructField("col1", IntegerType),
    StructField("col2", IntegerType),
    StructField("col3", IntegerType)
))

val finalDF = splitDF.select(from_json(col("json_col"), jsonSchema).as("data"))
    .select("data.*")

// 查看结果
finalDF.show()

方法二:自定义UDF提取字段

如果正则拆分的方式对你不够直观,可以用自定义UDF结合模式匹配提取各字段:

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

val rawDF = spark.read.text("your_file_path.csv")

// 定义匹配整行的正则:前三个是数字,最后是大括号包裹的JSON
val linePattern = "^(\\d+),(\\d+),(\\d+),({.*})$".r

// 自定义UDF提取字段
val extractData = udf((line: String) => line match {
    case linePattern(c1, c2, c3, jsonStr) => Some((c1.toInt, c2.toInt, c3.toInt, jsonStr))
    case _ => None
})

// 提取并过滤无效行
val parsedDF = rawDF.withColumn("extracted", extractData(col("value")))
    .filter(col("extracted").isNotNull)
    .select(
        col("extracted._1").as("c1"),
        col("extracted._2").as("c2"),
        col("extracted._3").as("c3"),
        col("extracted._4").as("json_col")
    )

// 解析JSON并得到最终结果
val jsonSchema = StructType(Seq(
    StructField("col1", IntegerType),
    StructField("col2", IntegerType),
    StructField("col3", IntegerType)
))

val finalDF = parsedDF.select(from_json(col("json_col"), jsonSchema).as("data"))
    .select("data.*")

finalDF.show()

两种方法最终都会生成你需要的DataFrame:

+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   2|   3|
|   2|   2|   4|
|   3|   3|   3|
+----+----+----+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 17:35:22