Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- import java.util.Arrays;
- import org.apache.spark.SparkConf;
- import org.apache.spark.api.java.JavaPairRDD;
- import org.apache.spark.api.java.JavaRDD;
- import org.apache.spark.api.java.JavaSparkContext;
- import scala.Tuple2;
- public class SparkWordCount {
- public static void main(String[] args) {
- System.out.println("Hello..sparkk.....");
- SparkConf conf = new SparkConf().setAppName("wordCountUsingSparkApis");
- JavaSparkContext sc = new JavaSparkContext(conf);
- sc.hadoopConfiguration().set("io.compression.codecs", "com.hadoop.compression.lzo.LzopCodec");
- JavaRDD<String> lines = sc.textFile("hdfs://krios/tmp/aree");
- JavaRDD<String> words =
- lines.flatMap(line -> Arrays.asList(line.split("\t")).iterator());
- JavaPairRDD<String, Integer> countData = words.filter(word -> !word.isEmpty()).mapToPair(t -> new Tuple2<>(t, 1)).reduceByKey((x, y) ->
- x + (int) y);
- countData.saveAsTextFile("hdfs://krios/tmp/SparkWordCountOutput");
- }
- }
Add Comment
Please, Sign In to add comment