Scala中含嵌套对象数组的DataFrame分组聚合最优实现方法
在Scala Spark DataFrame中处理对象数组的分组聚合
实现步骤与代码
要完成按前3列及数组元素的tag分组、对count求和并重组结构的需求,可通过以下步骤实现:
1. 展开嵌套数组
使用explode函数将_4.data数组拆分为独立行,让每个数组元素能单独参与计算。
2. 提取字段并分组求和
提取数组元素中的count和tag字段,以_1、_2、_3、tag为分组键,对count求和。
3. 聚合为目标数组结构
按_1、_2、_3重新分组,将每个分组的求和结果与tag组合成结构体,再收集为数组。
4. 重命名列名
将原列名修改为期望的UID、Title、Cell、Data。
完整Scala代码示例:
import org.apache.spark.sql.functions.{explode, sum, collect_list, struct} // 假设原始DataFrame为df val resultDF = df // 展开_4.data数组,每个元素转为单独一行 .select("_1", "_2", "_3", explode($"_4.data").alias("data_item")) // 提取结构体中的count和tag字段 .select( $"_1", $"_2", $"_3", $"data_item.count".alias("count"), $"data_item.tag".alias("tag") ) // 按指定维度分组,计算每个tag的总count .groupBy("_1", "_2", "_3", "tag") .agg(sum("count").alias("total_count")) // 重新聚合,将(total_count, tag)组合为结构体并收集成数组 .groupBy("_1", "_2", "_3") .agg(collect_list(struct($"total_count", $"tag")).alias("Data")) // 重命名列到目标名称 .withColumnRenamed("_1", "UID") .withColumnRenamed("_2", "Title") .withColumnRenamed("_3", "Cell")
代码说明
explode($"_4.data"):将嵌套数组拆分为多行,每行对应一个数组元素,解决嵌套结构无法直接分组的问题。- 两次
groupBy:第一次完成同tag的count求和,第二次将同一_1/_2/_3分组下的结果聚合成数组结构。 struct($"total_count", $"tag"):将求和结果与tag组合为结构体,保证输出数组的元素格式符合需求。
内容的提问来源于stack exchange,提问作者user3302967
相关产品推荐
相关产品推荐

