1、利用scala语言开发spark的worcount程序(本地运行)
1package com.zy.spark 2 3import org.apache.spark.rdd.RDD 4import org.apache.spark.{SparkConf, SparkContext} 5 6//todo:利用scala语言来实现spark的wordcount程序 7object WordCount { 8 def main(args: Array[String]): Unit = { 9 //1、创建SparkConf对象,设置appName和master local[2]表示本地采用2个线程去运行任务 10 val sparkConf: SparkConf = new SparkConf().setAppName("WordCount").setMaster("local[2]") 11 12 //2、创建SparkContext 该对象是所有spark程序的执行入口,它会创建DAGScheduler和TaskScheduler 13 val sc = new SparkContext(sparkConf) 14 15 //设置日志输出级别 16 sc.setLogLevel("warn") 17 18 //3、读取数据文件 19 val data: RDD[String] = sc.textFile("D:\\words.txt") 20 21 //4、切分每一行获取所有单词 22 val words: RDD[String] = data.flatMap(_.split(" ")) 23 24 //5、每个单词计为1 25 val wordAndOne: RDD[(String, Int)] = words.map((_, 1)) 26 27 //6、相同单词出现的所有的1累加 28 val result: RDD[(String, Int)] = wordAndOne.reduceByKey(_ + _) 29 30 //按照单词出现的次数降序排列 31 val sortRDD: RDD[(String, Int)] = result.sortBy(x => x._2, false) 32 33 34 //7、收集数据,打印输出 35 val finalResult: Array[(String, Int)] = sortRDD.collect() 36 finalResult.foreach(println) 37 38 //8、关闭sc 39 sc.stop() 40 } 41}
2、利用scala语言开发spark的wordcount程序(集群运行)
1package com.zy.spark 2 3import org.apache.spark.{SparkConf, SparkContext} 4import org.apache.spark.rdd.RDD 5 6//todo:利用scala语言开发spark的wordcount程序(集群运行) 7object WordCount_Online { 8 def main(args: Array[String]): Unit = { 9 //1、创建SparkConf对象,设置appName 10 val sparkConf: SparkConf = new SparkConf().setAppName("WordCount_Online") 11 12 //2、创建SparkContext 该对象是所有spark程序的执行入口,它会创建DAGScheduler和TaskScheduler 13 val sc = new SparkContext(sparkConf) 14 15 //设置日志输出级别 16 sc.setLogLevel("warn") 17 18 //3、读取数据文件 args(0)为文件地址参数 19 val data: RDD[String] = sc.textFile(args(0)) 20 21 //4、切分每一行获取所有单词 22 val words: RDD[String] = data.flatMap(_.split(" ")) 23 24 //5、每个单词计为1 25 val wordAndOne: RDD[(String, Int)] = words.map((_, 1)) 26 27 //6、相同单词出现的所有的1累加 28 val result: RDD[(String, Int)] = wordAndOne.reduceByKey(_ + _) 29 30 //7、把结果数据保存到hdfs上 args(1)是保存到hdfs的目录参数 31 result.saveAsTextFile(args(1)) 32 33 //8、关闭sc 34 sc.stop() 35 } 36 37}
最后打成jar包 到集群上执行
spark-submit --master spark://node1:7077 --class cn.itcast.spark.WordCount_Online --executor-memory 1g --total-executor-cores 2 original-spark_xxx-1.0-SNAPSHOT.jar /words.txt /out
3、利用java语言开发spark的wordcount程序(本地运行)
1package com.zy.spark; 2 3import org.apache.spark.SparkConf; 4import org.apache.spark.api.java.JavaPairRDD; 5import org.apache.spark.api.java.JavaRDD; 6import org.apache.spark.api.java.JavaSparkContext; 7import org.apache.spark.api.java.function.FlatMapFunction; 8import org.apache.spark.api.java.function.Function2; 9import org.apache.spark.api.java.function.PairFunction; 10import scala.Tuple2; 11 12import java.util.Arrays; 13import java.util.Iterator; 14import java.util.List; 15 16//todo:利用java语言开发spark的wordcount程序(本地运行) 17public class WordCount_Java { 18 public static void main(String[] args) { 19 //1、创建SparkConf对象 20 SparkConf sparkConf = new SparkConf().setAppName("WordCount_Java").setMaster("local[2]"); 21 22 //2、创建JavaSparkContext对象 23 JavaSparkContext jsc = new JavaSparkContext(sparkConf); 24 25 //3、读取数据文件 26 JavaRDD<String> data = jsc.textFile("D:\\words.txt"); 27 28 //4、切分每一行获取所有的单词 29 JavaRDD<String> words = data.flatMap(new FlatMapFunction<String, String>() { 30 public Iterator<String> call(String line) throws Exception { 31 String[] words = line.split(" "); 32 return Arrays.asList(words).iterator(); 33 } 34 }); 35 36 //5、每个单词计为1 37 JavaPairRDD<String, Integer> wordAndOne = words.mapToPair(new PairFunction<String, String, Integer>() { 38 public Tuple2<String, Integer> call(String word) throws Exception { 39 return new Tuple2<String, Integer>(word, 1); 40 } 41 }); 42 43 //6、相同单词出现1累加 44 JavaPairRDD<String, Integer> result = wordAndOne.reduceByKey(new Function2<Integer, Integer, Integer>() { 45 public Integer call(Integer v1, Integer v2) throws Exception { 46 return v1 + v2; 47 } 48 }); 49 50 //按照单词出现的次数降序排列 (单词,次数)------>(次数,单词).sortByKey------->(单词,次数) 51 52 JavaPairRDD<Integer, String> reverseRDD = result.mapToPair(new PairFunction<Tuple2<String, Integer>, Integer, String>() { 53 public Tuple2<Integer, String> call(Tuple2<String, Integer> t) throws Exception { 54 return new Tuple2<Integer, String>(t._2, t._1); 55 } 56 }); 57 58 JavaPairRDD<String, Integer> sortedRDD = reverseRDD.sortByKey(false).mapToPair(new PairFunction<Tuple2<Integer, String>, String, Integer>() { 59 public Tuple2<String, Integer> call(Tuple2<Integer, String> t) throws Exception { 60 return new Tuple2<String, Integer>(t._2, t._1); 61 } 62 }); 63 64 65 //7、收集数据打印输出 66 List<Tuple2<String, Integer>> finalResult = sortedRDD.collect(); 67 for (Tuple2<String, Integer> tuple : finalResult) { 68 System.out.println("单词:" + tuple._1 + " 次数:" + tuple._2); 69 } 70 71 //8、关闭jsc 72 jsc.stop(); 73 } 74}