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
相关产品推荐
相关产品推荐

