ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

SparkSQL 之 Json 格式数据转 DataSet 代码实现

2026/8/15 11:58:06 拓冰建站 浏览量
SparkSQL 之 Json 格式数据转 DataSet 代码实现 摘要JSON 是数据工程中最常见的半结构化格式如何高效地将 JSON 转为类型安全的 Dataset[CaseClass]本文从 spark.read.json() 的三种数据源、Schema 推断 vs 手动指定的优劣对比、嵌套 JSON 的三种展平模式Dot Notation / explode / from_json、以及 Encoder[CaseClass] 映射四个维度配合 2 张架构图 完整代码实例覆盖 Json Options 全部配置项给出生产级的 JSON→Dataset 转换实践。关键词spark.read.json, Schema 推断, explode, from_json, Encoder, CaseClass, Nested JSON, Json Options一、开篇JSON→Dataset 的正确姿势JSON → Dataset 三大核心问题 ① Schema: 自动推断 vs 手动指定 → 性能精确度权衡 ② Nested: Struct/Array/Map 嵌套字段如何展平 ③ Encoding: Row → CaseClass 的类型安全映射二、JSON → DataFrame/Dataset 转换流程2.1 三种数据源入口// ① 文件系统最常用valdfspark.read.json(hdfs://data/events/2024/*.json)valdfspark.read.json(/local/path/file.json)// ② RDD[String] → DataFramevaljsonStrings:Dataset[String]spark.createDataset(Seq({id:1,name:张三},{id:2,name:李四}))valdfspark.read.json(jsonStrings)// ③ DataFrame 直接构建valschemaStructType(Seq(StructField(id,LongType),StructField(name,StringType)))valdfspark.createDataFrame(rows,schema)2.2 Schema 推断 vs 手动指定// 方式 A: 自动推断方便但慢valdfspark.read.option(inferSchema,true).option(samplingRatio,0.1).json(path)// 方式 B: 手动 Schema推荐valschemaStructType(Seq(StructField(id,LongType),StructField(name,StringType),StructField(age,IntegerType)))valdfspark.read.schema(schema).json(path)// ✅ 零推断开销 · 类型精确 · 不依赖采样三、嵌套 JSON 展平三大模式3.1 Dot Notation — Struct 字段访问valflatdf.select($id,$name,$address.city.as(city),$address.street.as(street))// 或用 selectExpr SQL 风格valflatdf.selectExpr(id,name,address.city as city,address.street as street)3.2 explode — Array 数组展平1行→N行valexplodeddf.select($id,$name,explode($orders).as(order)).select($id,$name,$order.oid,$order.price)// explode_outer: 保留空数组行// posexplode: 额外输出数组下标3.3 from_json — 动态解析 JSON 字符串importorg.apache.spark.sql.functions.from_jsonvalorderSchemaStructType(Seq(StructField(oid,LongType),StructField(price,DoubleType)))df.select($id,from_json($jsonStrCol,orderSchema).as(parsed)).select($id,$parsed.oid,$parsed.price)四、Dataset[CaseClass] 类型安全映射caseclassUser(id:Long,name:String,age:Int)caseclassFlatOrder(id:Long,name:String,oid:Long,price:Double)valds:Dataset[FlatOrder]df.select($id,$name,explode($orders).as(order)).select($id,$name,$order.oid,$order.price).as[FlatOrder]// Encoder 自动推导// 类型安全操作ds.filter(_.price100).map(oo.copy(priceo.price*1.1))五、Json Options 速查类别参数说明损坏处理modePERMISSIVE/DROPMALFORMED/FAILFAST损坏行处理策略格式multiLine / allowComments / allowSingleQuotes非标准 JSON 兼容类型primitivesAsString / preferDecimal推断控制日期dateFormat / timestampFormat日期解析性能inferSchema / samplingRatio推断开销控制六、总结Schema 策略手动 Schema 比自动推断更快更精确生产环境推荐.schema(structType)。嵌套展平三模式Dot Notation(Struct) → explode(Array) → from_json(String JSON)。Dataset[CaseClass]通过.as[T]获得编译时类型安全Encoder 比 Kryo 快 10x。作者starzy博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践