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

如何以类似RDD的方式合并两个Dataset并聚合值为数组?

Spark DataFrame按Key聚合值为数组的实现方法

你可以通过合并DataFrame后分组聚合的方式实现需求,这是Spark中处理这类场景的惯用方法:

步骤1:合并两个DataFrame

先将两个DataFrame通过union合并,把所有key对应的value整合到同一列中:

val a = Map(1 -> 2, 3 -> 5).toSeq.toDF("key", "value")
val b = Map(1 -> 4, 2 -> 5).toSeq.toDF("key", "value")

val combinedDF = a.union(b)

步骤2:按Key分组并聚合值为数组

使用groupBy按key分组,再用collect_list函数将同一key的所有value聚合为数组:

import org.apache.spark.sql.functions.collect_list

val resultDF = combinedDF.groupBy("key").agg(collect_list("value").as("value"))
resultDF.show(false)

执行后输出结果:

+---+------+
|key|value |
+---+------+
|1  |[2, 4]|
|3  |[5]   |
|2  |[5]   |
+---+------+

仅保留两个DataFrame共有的Key

如果只需要两个DataFrame都存在的key(即内连接交集),可以先筛选共同key再聚合:

val commonKeysDF = a.select("key").intersect(b.select("key"))
val filteredResultDF = combinedDF.join(commonKeysDF, "key")
                                 .groupBy("key")
                                 .agg(collect_list("value").as("value"))
filteredResultDF.show(false)

输出结果:

+---+------+
|key|value |
+---+------+
|1  |[2, 4]|
+---+------+

替代方案:直接Join后合并列(仅适用于两个DataFrame场景)

如果仅处理两个DataFrame,也可以先做内连接,再将两列value合并为数组:

import org.apache.spark.sql.functions.array

a.join(b, "key")
  .select($"key", array($"value", $"value".alias("value_b")).as("value"))
  .show(false)

输出结果:

+---+------+
|key|value |
+---+------+
|1  |[2, 4]|
+---+------+

注:RDD的join实际返回(key, (v1, v2))格式的元组,类似cogroup才会返回包含迭代器的结构;DataFrame中推荐使用union+分组聚合的方式,该方法更通用,支持扩展到多个DataFrame的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:20:25