如何对Spark DataFrame分组合并数组、扁平化并提取唯一项?
解决方案:按ACC分组合并、扁平化并去重CCY数组
要实现你的需求,我们可以通过Spark的高阶函数和分组聚合操作来完成,下面分步骤拆解并给出两种可行的实现方式:
核心思路拆解
- 按ACC分组收集数组:先把同一ACC下的所有CCY数组收集到一起,形成嵌套数组(比如ACC=2会得到
[[AUD,YEN,USD], [GBP,AUD,YEN]]); - 扁平化+去重:把嵌套数组展开成一维数组,同时剔除重复的币种值。
方法一:flatten + array_distinct(直观简洁)
这种方式先把所有嵌套数组完全展开,再统一去重,代码逻辑一目了然,适合大多数场景:
import org.apache.spark.sql.functions.{collect_list, flatten, array_distinct} val resultDF = df .groupBy("ACC") .agg(collect_list("CCY").alias("nested_ccy")) // 收集同一ACC下的所有CCY数组 .select( $"ACC", array_distinct(flatten($"nested_ccy")).alias("unique_ccy") // 先扁平化再去重 )
方法二:aggregate + array_union(高效合并去重)
如果每组的数组数量多、元素量大,这种方式会更高效——它会逐个合并数组并即时去重,避免生成过大的中间数组:
import org.apache.spark.sql.functions.{collect_list, aggregate, lit, array_union} val resultDF = df .groupBy("ACC") .agg(collect_list("CCY").alias("nested_ccy")) .select( $"ACC", aggregate( $"nested_ccy", lit(Array.empty[String]).cast("array<string>"), // 初始值为空数组 (acc, arr) => array_union(acc, arr) // 逐个合并数组并去重 ).alias("unique_ccy") )
最终结果验证
两种方法运行后都会得到如下输出:
+---+-------------------------------+ |ACC|unique_ccy | +---+-------------------------------+ |1 |[USD, CAD] | |2 |[AUD, YEN, USD, GBP] | |3 |[BRL, AUS, BND, BOB, BWP, XAF, CLP, CNY, KMF, CSK, EGP]| +---+-------------------------------+
版本说明
上述用到的flatten、array_distinct、aggregate等高阶函数,需要Spark 2.4及以上版本支持。
内容的提问来源于stack exchange,提问作者1pluszara
相关产品推荐
相关产品推荐

