WordCount——自定义累加器实现WordCount

具体代码和步骤如下

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
}
经验分享 程序员 微信小程序 职场和发展