如何在Spark DataFrame(Scala API)中将类数组字符串转为数组?
优雅解决Spark DataFrame中类数组JSON字符串拆分多行的问题
没问题,针对你这种情况,用Spark内置的JSON处理函数就能优雅解决,完全不用手动去拆分字符串(手动拆分很容易踩坑,比如JSON里的逗号会干扰正常拆分逻辑)。下面是具体步骤和代码实现:
核心思路
- 定义目标JSON Schema:先明确
data列里的JSON数组对应的Spark Schema,让from_json能精准解析字符串。 - 解析字符串为数组类型:用
from_json把类数组格式的字符串转成Spark原生的数组结构体。 - 炸开数组为多行:用
explode函数把数组中的每个元素拆成单独的行。 - 可选:提取结构体字段:把解析后的结构体里的
@id、category等字段单独拆成列,方便后续业务处理。
完整Scala代码示例
import org.apache.spark.sql.functions.{from_json, explode, col} import org.apache.spark.sql.types.{StructType, StructField, StringType, ArrayType} // 你的原始DataFrame val df = Seq(("2019-05-15T10:37:22+00:00", """[{"@id":"1","@type":"type","category":"cat"},{"@id":"2","@type":"type","category":"cat1"}]""")).toDF("published", "data") // 定义data列对应的JSON Schema val dataSchema = ArrayType( StructType(Seq( StructField("@id", StringType, nullable = true), StructField("@type", StringType, nullable = true), StructField("category", StringType, nullable = true) )) ) // 解析JSON字符串 → 炸开数组 → 提取目标字段 val resultDF = df .withColumn("data_array", from_json(col("data"), dataSchema)) // 将字符串解析为数组结构体 .withColumn("data_item", explode(col("data_array"))) // 把数组拆分成单独行 .select( col("published"), col("data_item.@id").alias("id"), col("data_item.@type").alias("type"), col("data_item.category").alias("category") ) // 查看最终结果 resultDF.show(false)
运行结果
+------------------------+---+----+--------+ |published |id |type|category| +------------------------+---+----+--------+ |2019-05-15T10:37:22+00:00|1 |type|cat | |2019-05-15T10:37:22+00:00|2 |type|cat1 | +------------------------+---+----+--------+
为什么这是优雅的方案?
- 依赖Spark原生JSON解析能力,避免了手动处理字符串的复杂逻辑(比如转义字符、嵌套结构),稳定性更高。
- 类型安全:通过显式定义Schema,确保解析出来的结构符合预期,减少运行时错误。
- 代码简洁易读,完全贴合Spark的API风格,后续维护和扩展都很方便。
内容的提问来源于stack exchange,提问作者Evgeniy
相关产品推荐
相关产品推荐

