Spark Dataset分组聚合问题:expend列求和报错及列重复疑问
你的实现方式不正确,问题出在这两个核心点上
先看你的代码和目标SQL的明显差异:你写的SQL是按col1和col2分组求和,但代码里的groupBy额外加了expend列,这直接导致了两个问题:
1. 分组逻辑完全偏离预期
当你把expend加入groupBy时,Spark会把每个不同的expend值都当作独立的分组维度,相当于你是按col1 + col2 + expend的组合来分组。此时对每个分组求和expend,结果其实就是该分组本身的expend值——这完全不是你想要的“按col1、col2汇总expend”的效果。
2. 重复列的根源
groupBy中的列会被默认保留在结果集中,加上agg(sum("expend"))默认会生成一个名为expend的列,所以最终结果就出现了两个expend列:一个是分组时保留的原始列,另一个是求和后的列。
正确的实现方式
你需要对齐SQL的逻辑,只按col1和col2分组,同时给求和结果重命名避免列名冲突:
与你原代码结构对齐的写法
dataset.select(col("col1"), col("col2"), col("expend")) .groupBy(col("col1"), col("col2")) // 去掉expend分组 .agg(sum("expend").alias("total_expend")) // 给求和结果起别名
更简洁的写法(无需提前select)
dataset.groupBy("col1", "col2") .agg(sum("expend").alias("total_expend"))
这样执行后,结果的列会是[col1, col2, total_expend],完全匹配你目标SQL的输出,也不会有重复列的问题。
内容的提问来源于stack exchange,提问作者John Humanyun
相关产品推荐
相关产品推荐

