Scala中reduceGroups函数未达预期输出问题求助
问题:Spark中reduceGroups累加特定条件数值时mTotal始终为0
我现在尝试根据条件对数值求和,其中total变量需要在所有条件下累加,mTotal只在性别为男性("m")时累加。我写了下面的Spark代码,但运行后total能得到正确结果,mTotal却始终是0,麻烦帮忙看看哪里出问题了:
val record = file.map(rec => (rec.state,rec.gender,rec.Generated.toInt)) .groupByKey(_._1) .reduceGroups((a,b)=>{ var total:Int = 0 var mTotal:Int = 0 if(a._2.trim().equalsIgnoreCase("m")){ mTotal = a._3 + b._3 total = a._3 + b._3 }else{ total = a._3 + b._3 } (a._1,mTotal.toString(),total) }).collect
问题分析
你的代码逻辑存在几个关键问题,导致mTotal始终为0:
- 累加变量每次被重置:在
reduceGroups的匿名函数里,你每次都把total和mTotal初始化为0,完全丢弃了之前的累加结果——这意味着你每次只计算了当前两个元素的临时值,而不是持续累加所有元素。 - 混淆了累加结果与原始元素:
reduceGroups中,参数a是之前的累加汇总结果(也就是你返回的三元组),而b是当前遍历到的原始数据项。你现在把a当成了原始元素来处理,完全没用到它已经计算好的mTotal和total。 - 条件覆盖不完整:你只判断了
a的性别是否为男性,就算b是男性,也没有把b的数值加到mTotal里,逻辑完全漏掉了这种情况。
修正后的代码
我们需要调整逻辑,复用之前的累加结果,同时正确判断当前元素的性别来更新mTotal:
val record = file.map(rec => (rec.state, rec.gender, rec.Generated.toInt)) .groupByKey(_._1) .reduceGroups((acc, curr) => { // acc是历史累加结果:(state, 之前的mTotal字符串, 之前的total) // curr是当前要处理的单个数据项:(state, gender, Generated数值) val prevTotal = acc._3 val prevMTotal = acc._2.toInt // 将之前存储的字符串转回Int // 计算新的total:直接累加当前元素的数值 val newTotal = prevTotal + curr._3 // 计算新的mTotal:仅当当前元素是男性时累加,否则保留历史值 val newMTotal = if (curr._2.trim().equalsIgnoreCase("m")) { prevMTotal + curr._3 } else { prevMTotal } // 返回新的累加结果,保持结构一致 (acc._1, newMTotal.toString, newTotal) }).collect
更清晰的优化写法(推荐)
如果用case class来定义数据结构,代码可读性和可维护性会更高,也不容易出错:
// 定义case class明确数据结构 case class RawData(gender: String, value: Int) case class StateSummary(state: String, maleTotal: Int, total: Int) val record = file.map(rec => (rec.state, RawData(rec.gender, rec.Generated.toInt))) .groupByKey(_._1) .reduceGroups((summary, data) => { val newTotal = summary.total + data.value val newMaleTotal = if (data.gender.trim().equalsIgnoreCase("m")) { summary.maleTotal + data.value } else { summary.maleTotal } StateSummary(summary.state, newMaleTotal, newTotal) }) .map(_._2) // 只保留汇总结果 .collect
这种写法用类明确了每个变量的含义,避免了元组下标带来的混淆,后期修改逻辑也更方便。
内容的提问来源于stack exchange,提问作者Kumar Harsh
相关产品推荐
相关产品推荐

