Spark数据集扁平化:将数组列转为JSON序列化字符串列
解决Spark Dataset嵌套Schema扁平化时保留数组为JSON字符串的问题
我完全懂你的需求——在扁平化嵌套Schema的过程中,不想把数组字段展开成多行,而是直接把整个数组序列化成JSON字符串保留下来。用to_json确实是最贴合的方案,我来给你梳理完整的实现步骤和代码示例:
核心思路
- 遍历Dataset的
StructTypeSchema,区分普通嵌套字段和数组类型字段 - 对于普通嵌套字段:递归扁平化,用点号拼接完整字段路径(比如
user.address.city) - 对于数组类型字段:直接用
to_json函数将数组序列化为JSON字符串,用原路径作为别名
完整代码实现
1. 编写Schema遍历工具函数
先写一个递归函数,自动生成所有需要选择的列(包含扁平化的普通字段和转JSON的数组字段):
import org.apache.spark.sql.{Column, DataFrame} import org.apache.spark.sql.functions.{col, to_json} import org.apache.spark.sql.types.{ArrayType, StructType} def flattenSchemaWithArrayJson(df: DataFrame, parentPath: String = ""): Array[Column] = { df.schema.fields.flatMap { field => val currentPath = if (parentPath.isEmpty) field.name else s"$parentPath.${field.name}" field.dataType match { // 处理嵌套结构体,递归扁平化 case structType: StructType => flattenSchemaWithArrayJson(df.select(col(currentPath).as(field.name)), field.name) // 处理数组类型,直接转成JSON字符串 case arrayType: ArrayType => Array(to_json(col(currentPath)).alias(currentPath)) // 普通字段直接保留原路径别名 case _ => Array(col(currentPath).alias(currentPath)) } } }
2. 应用到你的Dataset
假设你的原始Dataset是originalDs,直接调用工具函数生成列列表,执行select即可完成扁平化:
val flattenedDs = originalDs.select(flattenSchemaWithArrayJson(originalDs): _*)
关键细节说明
- 字段命名自定义:当前代码用点号拼接嵌套路径作为别名,你可以根据需求修改(比如把点号换成下划线
currentPath.replace(".", "_")) - JSON序列化配置:如果需要控制JSON的生成规则(比如忽略null值、自定义日期格式),可以给
to_json传入配置参数,示例:to_json(col(currentPath), Map("ignoreNullFields" -> "true", "dateFormat" -> "yyyy-MM-dd")).alias(currentPath) - 深层嵌套兼容:函数会自动处理多层嵌套的结构体,不管嵌套层级有多深,都能正确扁平化;数组字段无论处于哪个嵌套层级,都会被转成JSON字符串保留
示例验证
假设你的原始Schema是这样的:
root |-- id: integer (nullable = false) |-- user: struct (nullable = true) | |-- name: string (nullable = true) | |-- hobbies: array (nullable = true) | | |-- element: string (containsNull = true) |-- orders: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- order_id: string (nullable = true) | | |-- amount: double (nullable = true)
扁平化后得到的Schema会是:
root |-- id: integer (nullable = false) |-- user.name: string (nullable = true) |-- user.hobbies: string (nullable = true) // 值为JSON字符串,例如["reading", "hiking"] |-- orders: string (nullable = true) // 值为JSON字符串,例如[{"order_id":"O1", "amount":100}, ...]
这样就完美实现了你想要的效果:普通嵌套字段彻底扁平化,数组字段保留为完整的JSON字符串而不被展开。
内容的提问来源于stack exchange,提问作者Kyle
相关产品推荐
相关产品推荐

