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

Scala与Spark环境下Delta表JSON字符串列的存储与解析问题

Delta表JSON字符串丢失键信息的解析问题

待解决问题

需要从Delta表中提取JSON字符串并解析,show函数可查看数据,但需将其转换为Map或样例类以进行处理。数据从JSON文件插入Delta表,表中对应列类型为String,但插入过程中仅存储了JSON数组的值,未保留键。执行select查询返回的是“不完整”数据,这是Delta表的默认行为吗?
示例:
待存储数据:{"name" : "Sample"}
实际存储内容:{"Sample"}

背景

已参考Stack Overflow相关问题并执行对应操作,但未找到解决方案。

插入步骤

数据从JSON文件插入Delta表,其中一个字段为JSON数组,需原样插入:

val jsonData = spark.read
  .option("multiLine", true)
  .option("mode", "PERMISSIVE")
  .option("dropFieldIfAllNull", false)
  .json(FILE_PATH)
  .createOrReplaceTempView("data")

val data = spark.sql("""SELECT * from Data""")

val inputDataFrame = data
  .select(
    col("Id"),
    col("Version"),
    col("StartDate"),
    col("EndDate"),
    explode(col("Configuration")).alias("Config")
  )
  .withColumn("ConfigId", col("Config.Id"))
  // 省略中间处理步骤
  .toDF
  
val deltaTable = DeltaTable.forPath(.....)

deltaTable
.as("Config") 
.merge(
       inputDataFrame.as("input"),
       // 省略匹配条件
)
.whenMatched
.update(
  Map(
      "ConfigData" -> col("input.Config"),
  )
)
.whenNotMatched
.insert(
  Map(
    "ConfigData" -> col("input.Config"),
  )
)
.execute()

查询步骤

执行select查询后,数据仅显示插入的JSON字符串的值,无键信息。

读取步骤

方法1

val arrayOffStringSchema = ArrayType(StringType)

var configData = df.select("Config").as[String]
configData.show(false) // 显示Config列,但数据无键信息
                                     
var desiredData = configData.withColumn("Config",from_json(col("Config"), arrayOffStringSchema))
desiredData.show(false) // 显示null

方法2

val json_schema = spark.read
      .option("multiLine", true)
      .option("mode", "PERMISSIVE") 
      .option("dropFieldIfAllNull", false)
      .json(df.select("Config").as[String]).schema
var config = configData.withColumn("Config", from_json(col("Config"),json_schema))
config.show(false) // 显示Config列,但数据无键信息
config.printSchema() // 输出如下结构:

/*
root
 |-- Config: struct (nullable = true)
 |    |-- _corrupt_record: string (nullable = true)
*/

平台信息

  • Scala 2.12.18
  • Apache Spark 3.5.1

疑问

请问哪里出错了?是插入过程有误,还是有其他实现方式?

编辑1

数据存储时会自动推断Schema,但该Schema未存储在Delta表中。因此读取数据时无法确定Schema来源,需要将其转换为Map以从Delta表的JSON字符串列中提取字段。


解决方案

问题根源

插入流程出错:你直接将Spark的StructType类型列input.Config赋值给String类型的ConfigData列,Spark会默认调用结构体的toString()方法,将其转为仅保留值的格式(如{Sample}),丢失键名。这不是Delta表的默认行为,而是Spark隐式类型转换导致的问题。

正确插入方式

插入时需用to_json函数将Struct类型转为标准JSON字符串:

// 更新和插入逻辑修改为:
.whenMatched
.update(
  Map(
      "ConfigData" -> to_json(col("input.Config")),
  )
)
.whenNotMatched
.insert(
  Map(
    "ConfigData" -> to_json(col("input.Config")),
  )
)

读取时的处理方法

方法1:转换为Map类型

若无需固定Schema,可直接将JSON字符串转为Map[String, String]:

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

val configData = df.select(
  from_json(col("ConfigData"), MapType(StringType, StringType)).alias("ConfigMap")
)
// 提取指定字段示例
configData.select("ConfigMap.name").show(false)

方法2:用样例类解析

若有固定Schema,可定义样例类后解析:

case class Config(name: String)
import org.apache.spark.sql.Encoders

val configSchema = Encoders.product[Config].schema
val configData = df.select(
  from_json(col("ConfigData"), configSchema).alias("Config")
)
configData.select("Config.name").show(false)

方法3:直接提取单个字段

若仅需特定字段,用get_json_object更高效:

val configData = df.select(
  get_json_object(col("ConfigData"), "$.name").alias("name")
)
configData.show(false)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 03:47:06