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
相关产品推荐
相关产品推荐

