如何使用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
相关产品推荐
相关产品推荐

