具体代码和步骤如下
def main(args: Array[String]): Unit = {
//1.创建SparkConf并设置App名称
val sparkConf = new SparkConf().setMaster("local[*]").setAppName("name")
//2.创建SparkContext,该对象是提交Spark App的入口
val sc= new SparkContext(sparkConf)
//6.注册累加器
val accu: WordCountAccu = new WordCountAccu
sc.register(accu,"wordcount") //可以为他取个名字:wordcount,也可以不取
//3.读取文件数据
val lineRDD: RDD[String] = sc.textFile("D:\workspace\spark\sparkCore\src\input\word1") //要计算的文件
//4.压平操作
val wordRDD: RDD[String] = lineRDD.flatMap(_.split(" ")) //按空格切分
//7.使用累加器计算WordCount——向累加器传入数据
wordRDD.foreach(word=>accu.add(word))
//8.取出累加器中的数据
val wordToCountMap: mutable.HashMap[String, Int] = accu.value
//9.打印
wordToCountMap.foreach(println)
//10.关闭连接
sc.stop()
}
}
//5.自定义累加器 In Out
class WordCountAccu extends AccumulatorV2[String,mutable.HashMap[String,Int]]{
private val map = new mutable.HashMap[String,Int]()
//判空
override def isZero: Boolean = map.isEmpty
//复制
override def copy(): AccumulatorV2[String, mutable.HashMap[String, Int]] = new WordCountAccu
//重置
override def reset(): Unit = map.clear()
//区内累加数据
override def add(v: String): Unit = {
map(v) = map.getOrElse(v,0) + 1
}
//区间合并数据
override def merge(other: AccumulatorV2[String, mutable.HashMap[String, Int]]): Unit = {
other.value.foreach{
case (word,count) =>
map(word) = map.getOrElse(word,0) + count
}
}
//返回值
override def value: mutable.HashMap[String, Int] = map
}