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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 02:25:47