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

Spark Scala:按visitor分组并将对应数组聚合为单个数组

解决方案:按Visitor分组聚合Asset数组

没问题,这是Spark中处理数组聚合的常见场景,我来给你一步步实现需求:

1. 先修正原始数据的创建方式

你提供的示例数据写法有问题(Seq里混合了字符串和数组,会导致类型不匹配),正确的做法是把每个访客和对应的资产数组包装成元组:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{collect_list, flatten, array_distinct}

// 初始化SparkSession(如果还没初始化的话)
val spark = SparkSession.builder().appName("AssetAggregation").master("local[*]").getOrCreate()
import spark.implicits._

// 正确构造原始数据:每个元素是(visitor, asset数组)的元组
val rawData = Seq(
  ("visitor1", Array("item1","item2","item3","item4")),
  ("visitor2", Array("item1","item2","item3")),
  ("visitor1", Array("item1","item2","item5")),
  ("visitor2", Array("item4","item6"))
)

// 转换为DataFrame
val df = rawData.toDF("visitor", "asset")

2. 执行分组聚合

核心思路是:先按visitor分组,收集每组内的所有asset数组,再把嵌套的数组展开为一维数组:

// 基础聚合(保留所有元素,包括重复项)
val aggregatedDF = df.groupBy("visitor")
  .agg(flatten(collect_list("asset")).alias("aggregated_asset"))

// 查看结果
aggregatedDF.show(false)

执行后输出结果:

+--------+---------------------------------------+
|visitor |aggregated_asset                       |
+--------+---------------------------------------+
|visitor1|[item1, item2, item3, item4, item1, item2, item5]|
|visitor2|[item1, item2, item3, item4, item6]     |
+--------+---------------------------------------+

可选:去重聚合

如果需要去掉资产数组中的重复元素,可以在聚合时加上array_distinct:

// 去重后的聚合
val aggregatedDistinctDF = df.groupBy("visitor")
  .agg(array_distinct(flatten(collect_list("asset"))).alias("aggregated_asset_distinct"))

aggregatedDistinctDF.show(false)

输出结果(重复的item被移除):

+--------+-----------------------------------+
|visitor |aggregated_asset_distinct          |
+--------+-----------------------------------+
|visitor1|[item1, item2, item3, item4, item5]|
|visitor2|[item1, item2, item3, item4, item6]|
+--------+-----------------------------------+

关键函数解释

  • groupBy("visitor"):将数据按访客ID分组,把同一访客的所有行归为一组
  • collect_list("asset"):收集每组内的所有asset数组,得到一个数组的数组(比如visitor1会得到Array(Array("item1",...), Array("item1",...)))
  • flatten(...):将嵌套的数组展开为一维数组,把所有子数组的元素合并到一个数组中
  • array_distinct(...):可选函数,用于移除数组中的重复元素

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:22:10