Skip to content
中文
1 min read#spark

11 Common Word Count Methods

Common wordCount in Spark

Updated:

阅读中文版

import org.apache.spark.rdd.RDD import org.apache.spark.{SparkConf, SparkContext}

import scala.collection.mutable

/**

  • @Autor LZH

  • @Date 2020/12/1 10:17 / object WordCount extends App { var conf=new SparkConf().setMaster("local[]").setAppName("WordCount") val sc=new SparkContext(conf)

    val rdd: RDD[(String, Int)] = sc.makeRDD( List( ("a", 1), ("a", 2), ("c", 1), ("b", 1), ("b", 2), ("c", 2) ), 2 )

println("Method 1: reduceByKey()") rdd.reduceByKey(+).collect().foreach(println)

println("Method 2: groupByKey()") rdd.groupByKey().mapValues(_.sum).collect().foreach(println)

println("Method 3: groupBy()") rdd.groupBy(.1).mapValues(.map(._2).sum).collect().foreach(println)

println("Method 4: aggregateByKey()") rdd.aggregateByKey(0)(+,+).collect().foreach(println)

println("Method 5: foldByKey()") rdd.foldByKey(0)(+).collect().foreach(println)

println("Method 6: combineByKey()") rdd.combineByKey( v=>v, (x: Int, y) => x+y, (x: Int, y: Int) => x+y ).collect().foreach(println)

println("Method 7: countByKey()") rdd.map{ case (str,sum)=> (str+" ")*sum }.flatMap(.split(" ")).map((,1)).countByKey().foreach(println)

println("Method 8: countByValue()") rdd.map{ case (str,sum)=> (str+" ")*sum }.flatMap(_.split(" ")).countByValue().foreach(println)

println("Method 9: aggregate()") rdd.map{ case (str,sum)=> (str+" ")*sum }.flatMap(_.split(" ")).map(s => mutable.Map(s -> 1)) .aggregate(mutable.MapString, Int)( (map1: mutable.Map[String, Int], map2: mutable.Map[String, Int]) => { map1.foldLeft(map2)( (innerMap, kv) => { innerMap(kv._1) = innerMap.getOrElse(kv._1, 0) + kv._2 innerMap } ) }, (map1: mutable.Map[String, Int], map2: mutable.Map[String, Int]) => { map1.foldLeft(map2)( (innerMap, kv) => { innerMap(kv._1) = innerMap.getOrElse(kv._1, 0) + kv._2 innerMap } ) } ).foreach(println)

println("Method 10: fold()") rdd.map(t => mutable.Map((t._1,t._2))).fold(mutable.MapString, Int)( (map1, map2) => { map2.foreach{ case (word,cnt)=>{ val oldCnt= map1.getOrElse(word,0) map1.update(word,oldCnt+cnt) } } map1 } ).foreach(println)

println("Method 11: reduce()") rdd.map(t=>mutable.Map((t._1,t._2))).reduce( (map1, map2) => { map2.foreach{ case (word,cnt)=>{ val oldCnt= map1.getOrElse(word,0) map1.update(word,oldCnt+cnt) } } map1 } ).foreach(println)

}

Related posts

By shared tags

Comments(0)