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

Scala中如何从Spark推荐结果DataFrame生成两列新DataFrame?

解决Spark推荐结果DataFrame格式转换问题

嘿,看你正在用Spark做协同过滤推荐,已经拿到用户31511的Top10推荐结果啦,现在要把嵌套的数组结构拆成扁平的productionId和rating两列,这个不难,我给你一步步说:

首先,先回顾下你当前results的结构和输出:

当前results的输出和Schema

scala> results.show
+------+--------------------+
|userId| recommendations|
+------+--------------------+
| 31511|[[328, 0.7845393]...|
+------+--------------------+

scala> results.printSchema
root
 |-- userId: integer (nullable = false)
 |-- recommendations: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- productionId: integer (nullable = true)
 |    |    |-- rating: float (nullable = true)

你的目标是转换成这样的扁平结构:

+------------+------+
|productionId|rating|
+------------+------+
| 1          | 1    |
+------------+------+
| 2          | 4.5  |
+------------+------+
| 3          | 5    |
+------------+------+
| 4          | 2.5  |
+------------+------+

实现步骤

要把数组列拆成多行,我们需要用Spark的explode函数,它能把数组里的每个元素拆成单独的行。具体代码如下:

// 先导入explode函数
import org.apache.spark.sql.functions.explode

// 转换DataFrame
val flattenedDF = results
  .select(explode($"recommendations").alias("rec")) // 展开数组列,给每个struct元素起别名rec
  .select($"rec.productionId", $"rec.rating") // 从struct里提取需要的字段

// 查看结果
flattenedDF.show()

这样就能得到你想要的扁平结构啦!如果还需要保留userId字段,也可以在select的时候加上$"userId",比如:

val flattenedDFWithUserId = results
  .select($"userId", explode($"recommendations").alias("rec"))
  .select($"userId", $"rec.productionId", $"rec.rating")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:16:17