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

