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

如何对Spark DataFrame分组合并数组、扁平化并提取唯一项?

解决方案:按ACC分组合并、扁平化并去重CCY数组

要实现你的需求,我们可以通过Spark的高阶函数和分组聚合操作来完成,下面分步骤拆解并给出两种可行的实现方式:

核心思路拆解

  1. 按ACC分组收集数组:先把同一ACC下的所有CCY数组收集到一起,形成嵌套数组(比如ACC=2会得到[[AUD,YEN,USD], [GBP,AUD,YEN]]);
  2. 扁平化+去重:把嵌套数组展开成一维数组,同时剔除重复的币种值。

方法一: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 10:01:25