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

如何在Apache Spark Scala中实现DataFrame列映射与键值对归约

用Spark基础操作实现女演员-电影列表的按键归约

没问题!咱们用Spark的核心基础操作(不用Spark SQL)来实现这个需求,一步步给你讲清楚:


步骤1:转换为键值对格式数据

首先,从你的DataFrame里只选出需要的actress和movie title两列,然后把每一行映射成*(女演员名字, 电影标题)*的键值对元组。这一步是为了让数据适配Spark按键操作的格式要求。

Python 示例代码:

# 假设你的DataFrame名为df
key_value_rdd = df.select("actress", "movie title").rdd.map(lambda row: (row[0], row[1]))

Scala 示例代码:

// 假设你的DataFrame名为df
val keyValueRDD = df.select("actress", "movie title").rdd.map(row => (row.getString(0), row.getString(1)))

步骤2:按键归约,收集每位女演员的电影列表

接下来通过Spark的按键聚合操作,把同一个女演员对应的所有电影标题收集成一个列表。这里有两种常用实现方式,更推荐第二种(性能更优):

方式1:用groupByKey(简单直观)

groupByKey会直接按键分组,之后我们把分组后的迭代器转换为列表即可:

Python:
actress_movies_rdd = key_value_rdd.groupByKey().mapValues(list)
Scala:
val actressMoviesRDD = keyValueRDD.groupByKey().mapValues(_.toList)

方式2:用aggregateByKey(高效聚合,减少shuffle开销)

如果数据量较大,aggregateByKey比groupByKey性能更好——它会先在每个分区本地完成部分聚合,再进行全局shuffle,避免冗余数据传输:

Python:
# 初始化空列表,本地聚合时添加电影标题,全局聚合时合并列表
actress_movies_rdd = key_value_rdd.aggregateByKey([], lambda acc, movie: acc + [movie], lambda acc1, acc2: acc1 + acc2)
Scala:
// 初始化空列表,本地聚合和全局聚合分别完成列表拼接
val actressMoviesRDD = keyValueRDD.aggregateByKey(List.empty[String])(
  (acc, movie) => movie :: acc,
  (acc1, acc2) => acc1 ++ acc2
)

步骤3:(可选)转回DataFrame格式

如果需要把结果转回DataFrame方便后续分析或存储,只需调用toDF方法:

Python:

result_df = actress_movies_rdd.toDF("actress", "movie_list")

Scala:

val resultDF = actressMoviesRDD.toDF("actress", "movie_list")

这样你就得到了以女演员为键、对应电影列表为值的最终结果啦~

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:23:14