You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame如何将字符串数组列与event_id关联展开?

Spark DataFrame JSON数组字符串转换方案

问题背景

现有Spark DataFrame结构如下:

scala> df1.printSchema
root
 |-- actions: string (nullable = true)
 |-- event_id: string (nullable = true)

其中actions列存储的是JSON对象数组字符串,但类型为String,无法直接使用explode函数展开。

示例输入数据:

event_idactions
1[{"name": "Vijay", "score": 843},{"name": "Manish", "score": 840}, {"name": "Mayur", "score": 930}]

期望输出格式:

event_idnamescore
1Vijay843
1Manish840
1Mayur930

之前尝试直接读取actions列的JSON,但丢失了event_id关联:

val df2= spark.read.option("multiline",true).json(df1.rdd.map(row => row.getAs[String]("actions")))

解决方案

方法一:DataFrame API 解析(推荐)

利用Spark内置的from_json函数直接解析JSON字符串,全程使用DataFrame API,性能更优且无需转换到RDD。

import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._

// 定义单个JSON对象的Schema
val actionItemSchema = StructType(Seq(
  StructField("name", StringType, nullable = true),
  StructField("score", IntegerType, nullable = true)
))

// 执行转换
val resultDF = df1
  // 将actions字符串解析为JSON数组
  .withColumn("actions_parsed", from_json(col("actions"), ArrayType(actionItemSchema)))
  // 展开数组为多行
  .withColumn("action", explode(col("actions_parsed")))
  // 提取字段并保留event_id
  .select(
    col("event_id"),
    col("action.name").alias("name"),
    col("action.score").alias("score")
  )

// 查看结果
resultDF.show()

方法二:RDD转换(兼容旧版本)

如果需要用RDD处理,需保留event_id与actions的关联,避免丢失字段:

import org.apache.spark.sql.Row

// 保留event_id和actions的关联,转换为(event_id, actions_str)格式的RDD
val eventActionPairs = df1.rdd.map(row => (row.getAs[String]("event_id"), row.getAs[String]("actions")))

// 解析JSON数组并扁平化,生成(event_id, name, score)的RDD
val flattenedRDD = eventActionPairs.flatMap { case (eventId, actionsStr) =>
  // 解析单个actions字符串为DataFrame,再转成数组
  val actionRows = spark.read.json(spark.sparkContext.parallelize(Seq(actionsStr))).collect()
  // 遍历每个JSON对象,提取字段并关联event_id
  actionRows.map(row => (eventId, row.getAs[String]("name"), row.getAs[Int]("score")))
}

// 转换为目标DataFrame
val resultDF = flattenedRDD.toDF("event_id", "name", "score")

resultDF.show()

内容的提问来源于stack exchange,提问作者Manish Kumar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 02:55:43