Guest User

Untitled

a guest
Nov 28th, 2018
74
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
text 0.99 KB | None | 0 0
  1.  
  2. import java.util.Arrays;
  3. import org.apache.spark.SparkConf;
  4. import org.apache.spark.api.java.JavaPairRDD;
  5. import org.apache.spark.api.java.JavaRDD;
  6. import org.apache.spark.api.java.JavaSparkContext;
  7. import scala.Tuple2;
  8.  
  9. public class SparkWordCount {
  10.  
  11. public static void main(String[] args) {
  12. System.out.println("Hello..sparkk.....");
  13. SparkConf conf = new SparkConf().setAppName("wordCountUsingSparkApis");
  14. JavaSparkContext sc = new JavaSparkContext(conf);
  15.  
  16. sc.hadoopConfiguration().set("io.compression.codecs", "com.hadoop.compression.lzo.LzopCodec");
  17.  
  18. JavaRDD<String> lines = sc.textFile("hdfs://krios/tmp/aree");
  19. JavaRDD<String> words =
  20. lines.flatMap(line -> Arrays.asList(line.split("\t")).iterator());
  21.  
  22. JavaPairRDD<String, Integer> countData = words.filter(word -> !word.isEmpty()).mapToPair(t -> new Tuple2<>(t, 1)).reduceByKey((x, y) ->
  23. x + (int) y);
  24.  
  25. countData.saveAsTextFile("hdfs://krios/tmp/SparkWordCountOutput");
  26.  
  27. }
  28. }
Add Comment
Please, Sign In to add comment