Spark Java API中分组列与聚合的实现及代码修正问题
问题分析与解决方案
咱们来拆解你遇到的Spark聚合问题:你需要基于col1和col2的组合对数据集分组,对expend1和expend2求和,但最初的代码没有输出预期的聚合结果,写入CSV后数据不符合预期。
核心问题点
- 分组逻辑不完整:你尝试的示例代码中,
groupBy只指定了col1,但从原始数据和预期结果来看,你需要按col1和col2共同分组(比如itemA1 + itemA2是一组),缺失col2会导致分组错误,无法得到正确的聚合值。 - 聚合列未指定别名:最初的
agg(sum(col("expend1")), sum(col("expend2")))会生成默认列名(比如sum(expend1)),而你期望输出的列名还是expend1和expend2,这不仅会导致CSV表头不符合预期,还可能让你误以为聚合未生效。
修正后的正确代码
结合你的需求,同时修正分组逻辑和列别名问题的代码如下:
sampleDS = sampleDS.select(col("col1"), col("col2"), col("expend1"), col("expend2")) .groupBy(col("col1"), col("col2")) // 按两个列的组合分组 .agg( sum(col("expend1")).as("expend1"), // 给聚合结果指定目标列名 sum(col("expend2")).as("expend2") ) .sort(col("col1"), col("col2"));
你最终使用的expr("sum(expend1)").as("expend1")写法也是完全可行的,和上面的代码效果一致,都是为聚合结果指定别名。
验证聚合结果
修正后,数据集会按照col1和col2的组合正确分组,对支出列求和,得到你预期的结果:
+-------+-------+-------+-------+ | col1 | col2 |expend1|expend2| +-------+-------+-------+-------+ |itemA1 | itemA2| 10 | 11 | |itemB1 | itemB2| 9 | 7 | +-------+-------+-------+-------+
关于CSV写入
你的CSV写入代码是正确的:coalesce(1)将数据合并为单个文件,SaveMode.Overwrite覆盖已有文件,option("header", true)写入表头,配合修正后的聚合代码就能输出符合预期的CSV文件。
内容的提问来源于stack exchange,提问作者John Humanyun
相关产品推荐
相关产品推荐

