Spark算子之AggregateByKey使用及流程解析*
package com.bigdata.spark.core.rdd.oper.transform
import org.apache.spark.{SparkConf, SparkContext}
object RDD_Oper_Transform_1 {
def main(args: Array[String]): Unit = {
val conf = new SparkConf().setMaster("local[*]").setAppName("Transform")
val sc = new SparkContext(conf)
val rdd = sc.makeRDD(
List(
("a", 1), ("a", 2), ("b", 3),
("b", 4), ("b", 5), ("a", 6)
),
2
)
// aggregateByKey算子可以区分分区内和分区间的计算逻辑,就意味着分区内和分区间的计算逻辑可以不同。但是也可以相同
// TODO aggregateByKey算子可以实现 WordCount ( 4 / 10 )
// rdd.aggregateByKey(0)(
// (x, y) => {
// x + y
// },
// (x, y) => {
// x + y
// }
// ).collect().foreach(println)
//rdd.aggregateByKey(0)(_+_,_+_).collect().foreach(println)
// TODO 如果aggregateByKey算子在使用时,分区内计算规则和分区间计算规则相同,那么可以采用新的算子代替:foldByKey
// foldByKey算子可以实现 WordCount ( 5 / 10 )
rdd.foldByKey(0)(_+_).collect().foreach(println)
sc.stop()
}
}
AggregateByKey使用及流程解析



