Scala与Spark环境下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

