如何使用Spark Scala将字符串格式的复杂JSON转换为DataFrame
错误根因
你定义的schema要求Row包含Array[String]类型的info字段和String类型的more字段,但你传入Row的是完整的JSON字符串,Spark无法直接将整段字符串映射为你声明的复合结构,因此抛出类型不匹配异常。
实现方案
方案1:自动推导Schema(推荐,代码最简)
直接使用Spark原生的JSON读取能力解析字符串RDD,Spark会自动根据JSON结构推导Schema,无需手动定义,适配性最高。
import org.apache.spark.sql.SparkSession // 初始化SparkSession val ss = SparkSession.builder() .appName("JsonToDataFrame") .master("local[*]") // 本地测试用,生产环境移除该配置 .getOrCreate() // 你的原始JSON字符串 val data = """{ "info": [ { "done": "time", "id": 9, "type": "normal", "pid": 202020, "add": { "fields": true, "stat": "not sure" } }, { "done": "time", "id": 14, "type": "normal", "pid": 764310, "add": { "fields": true, "stat": "sure" } }, { "done": "time", "id": 9, "type": "normal", "pid": 202020, "add": { "note": { "id": 922, "score": 0 } } } ], "more": { "a": "ok", "b": "fine", "c": 3 } }""" // 构造字符串RDD后直接用json接口读取 val jsonRdd = ss.sparkContext.parallelize(Seq(data)) val df = ss.read.json(jsonRdd) // 验证结果 df.printSchema() df.show(false)
方案2:自定义Schema(适用于需要严格控制字段类型的场景)
如果需要避免自动推导的类型偏差,可以自定义和JSON结构完全匹配的Schema,通过from_json函数完成解析:
import org.apache.spark.sql.functions.{col, from_json} import org.apache.spark.sql.types._ // 逐层定义和JSON结构匹配的Schema val addSchema = new StructType() .add("fields", BooleanType) .add("stat", StringType) .add("note", new StructType() .add("id", LongType) .add("score", IntegerType) ) val infoItemSchema = new StructType() .add("done", StringType) .add("id", LongType) .add("type", StringType) .add("pid", LongType) .add("add", addSchema) val rootSchema = new StructType() .add("info", ArrayType(infoItemSchema)) .add("more", new StructType() .add("a", StringType) .add("b", StringType) .add("c", IntegerType) ) // 先将JSON字符串转为单列DataFrame,再解析 val rawDf = ss.createDataFrame(Seq((data,))).toDF("json_str") val parsedDf = rawDf.select(from_json(col("json_str"), rootSchema).alias("root")) val finalDf = parsedDf.select("root.*") // 验证结果 finalDf.printSchema() finalDf.show(false)
内容的提问来源于stack exchange,提问作者Arjun
相关产品推荐
相关产品推荐

