如何在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
相关产品推荐
相关产品推荐

