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

如何使用Spark Scala分组构建嵌套DataFrame并解决聚合报错?

解决Spark DataFrame嵌套Struct的报错问题

错误原因

你使用groupBy的方式有误:Spark要求groupBy后的聚合操作中,所有未被纳入分组的字段必须用聚合函数(比如first()、max()等)包裹。但你的需求只是将number和date字段打包成嵌套的struct,完全不需要分组操作——原数据中id、pseudo_id、system这三个字段的组合本身就是唯一的,没有重复行需要聚合。

正确实现代码

直接通过select语句重组字段即可,不需要groupBy:

val nestedDf: DataFrame = df
  .select(
    col("id"),
    col("pseudo_id"),
    col("system"),
    struct(
      col("number"),
      col("date")
    ).as("version")
  )

验证结果

执行上述代码后,nestedDf的schema会完全符合你的要求:

root
 |-- id: string (nullable = true)
 |-- pseudo_id: string (nullable = true)
 |-- system: string (nullable = true)
 |-- version: struct (nullable = false)
 |    |-- number: string (nullable = true)
 |    |-- date: string (nullable = true)

补充说明

如果你的数据中存在id、pseudo_id、system重复的行,确实需要分组聚合,那可以用first()聚合函数包裹字段,代码如下:

val nestedDf: DataFrame = df
  .groupBy("id", "pseudo_id", "system")
  .agg(
    struct(
      first("number").as("number"),
      first("date").as("date")
    ).as("version")
  )

这种写法会取每组中第一条数据的number和date来构建version struct,适合存在重复分组键的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 12:01:21