From 42691bef737b9a6470cecc0c71a45733805fe185 Mon Sep 17 00:00:00 2001 From: greg Date: Thu, 8 Feb 2018 19:52:49 +0300 Subject: [PATCH 1/5] added my wc --- .gitignore | 6 +++ src/main/java/greg/WordCount.java | 63 ++++++++++++++++++++++++++++++ src/main/java/greg/WordsOrder.java | 13 ++++++ 3 files changed, 82 insertions(+) create mode 100644 src/main/java/greg/WordCount.java create mode 100644 src/main/java/greg/WordsOrder.java diff --git a/.gitignore b/.gitignore index e3b6fa2..99de303 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,9 @@ .idea /input /output* +2017-big-data.iml +classes/ +src/main/java/META-INF/ +src/main/java/greg/META-INF/ +target/ + diff --git a/src/main/java/greg/WordCount.java b/src/main/java/greg/WordCount.java new file mode 100644 index 0000000..daad4d2 --- /dev/null +++ b/src/main/java/greg/WordCount.java @@ -0,0 +1,63 @@ +package greg; + +import java.io.IOException; +import java.util.StringTokenizer; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Job; +import org.apache.hadoop.mapreduce.Mapper; +import org.apache.hadoop.mapreduce.Reducer; +import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; +import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; + +public class WordCount { + + public static class TokenizerMapper + extends Mapper{ + + private final static IntWritable one = new IntWritable(1); + private Text word = new Text(); + + public void map(Object key, Text value, Context context + ) throws IOException, InterruptedException { + StringTokenizer itr = new StringTokenizer(value.toString()); + while (itr.hasMoreTokens()) { + word.set(itr.nextToken()); + context.write(word, one); + } + } + } + + public static class IntSumReducer + extends Reducer { + private IntWritable result = new IntWritable(); + + public void reduce(Text key, Iterable values, + Context context + ) throws IOException, InterruptedException { + int sum = 0; + for (IntWritable val : values) { + sum += val.get(); + } + result.set(sum); + context.write(key, result); + } + } + + public static void main(String[] args) throws Exception { + Configuration conf = new Configuration(); + Job job = Job.getInstance(conf, "word count"); + job.setJarByClass(WordCount.class); + job.setMapperClass(TokenizerMapper.class); + job.setCombinerClass(IntSumReducer.class); + job.setReducerClass(IntSumReducer.class); + job.setOutputKeyClass(Text.class); + job.setOutputValueClass(IntWritable.class); + FileInputFormat.addInputPath(job, new Path(args[0])); + FileOutputFormat.setOutputPath(job, new Path(args[1])); + System.exit(job.waitForCompletion(true) ? 0 : 1); + } +} diff --git a/src/main/java/greg/WordsOrder.java b/src/main/java/greg/WordsOrder.java new file mode 100644 index 0000000..f0772d6 --- /dev/null +++ b/src/main/java/greg/WordsOrder.java @@ -0,0 +1,13 @@ +package greg; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Mapper; + +public class WordsOrder { + + public static class MyMapper extends Mapper { + + } +} From f4e652bad8932f762c16dead493814cc4f901574 Mon Sep 17 00:00:00 2001 From: greg Date: Fri, 9 Feb 2018 01:44:59 +0300 Subject: [PATCH 2/5] sorting words complete --- src/main/java/greg/SortWords.java | 83 ++++++++++++++++++++++++++++++ src/main/java/greg/WordCount.java | 24 +++++++-- src/main/java/greg/WordsOrder.java | 13 ----- 3 files changed, 104 insertions(+), 16 deletions(-) create mode 100644 src/main/java/greg/SortWords.java delete mode 100644 src/main/java/greg/WordsOrder.java diff --git a/src/main/java/greg/SortWords.java b/src/main/java/greg/SortWords.java new file mode 100644 index 0000000..870b6d6 --- /dev/null +++ b/src/main/java/greg/SortWords.java @@ -0,0 +1,83 @@ +package greg; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.conf.Configured; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.io.WritableComparator; +import org.apache.hadoop.mapreduce.Job; +import org.apache.hadoop.mapreduce.Mapper; +import org.apache.hadoop.mapreduce.Reducer; +import org.apache.hadoop.mapreduce.lib.input.TextInputFormat; +import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; +import org.apache.hadoop.util.Tool; +import org.apache.hadoop.util.ToolRunner; +import java.io.IOException; + +public class SortWords extends Configured implements Tool { + + public static class MyMapper extends Mapper { + + @Override + protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { + final String[] data = value.toString().split("\t"); + context.write(new IntWritable(Integer.valueOf(data[1])), new Text(data[0])); + } + } + + public static class MyReducer extends Reducer { + + @Override + protected void reduce(IntWritable key, Iterable values, Context context) throws IOException, InterruptedException { + super.reduce(key, values, context); + } + } + + public static class Comparetor extends WritableComparator { + + @Override + public int compare(byte[] b1, int s1, int l1, byte[] b2, int s2, int l2) { + return -Integer.compare(readInt(b1, s1), readInt(b2, s2)); + } + } + + @Override + public int run(String[] args) throws Exception { + final Job job = new Job(new Configuration(), "sort words"); + job.setJarByClass(WordCount.class); + + job.setMapperClass(MyMapper.class); + job.setReducerClass(MyReducer.class); + job.setSortComparatorClass(Comparetor.class); + + TextInputFormat.addInputPath(job, new Path(args[1])); + TextOutputFormat.setOutputPath(job, new Path(args[1] + "_2")); + + job.setInputFormatClass(TextInputFormat.class); + //job.setInputFormatClass(SequenceFileInputFormat.class); + //job.setOutputFormatClass(TextOutputFormat.class); + job.setMapOutputKeyClass(IntWritable.class); + job.setMapOutputValueClass(Text.class); + + job.setOutputKeyClass(Text.class); + job.setOutputValueClass(IntWritable.class); + + return job.waitForCompletion(true) ? 0 : 1; + } + + public static void main(final String[] args) throws Exception { + + if(args.length != 2) { + throw new IllegalArgumentException("Usage: hadoop jar .jar "); + } + + final int returnCode1 = ToolRunner.run(new Configuration(), new greg.WordCount(), args); + final int returnCode2 = ToolRunner.run(new Configuration(), new greg.SortWords(), args); + + if(0 == returnCode1 && 0 == returnCode2) + System.exit(0); + System.exit(1); + } +} diff --git a/src/main/java/greg/WordCount.java b/src/main/java/greg/WordCount.java index daad4d2..ada8731 100644 --- a/src/main/java/greg/WordCount.java +++ b/src/main/java/greg/WordCount.java @@ -4,6 +4,7 @@ import java.util.StringTokenizer; import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.conf.Configured; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; @@ -12,8 +13,10 @@ import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; +import org.apache.hadoop.util.Tool; +import org.apache.hadoop.util.ToolRunner; -public class WordCount { +public class WordCount extends Configured implements Tool { public static class TokenizerMapper extends Mapper{ @@ -47,17 +50,32 @@ public void reduce(Text key, Iterable values, } } - public static void main(String[] args) throws Exception { + @Override + public int run(String[] args) throws Exception { + Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "word count"); + job.setJarByClass(WordCount.class); + job.setMapperClass(TokenizerMapper.class); job.setCombinerClass(IntSumReducer.class); job.setReducerClass(IntSumReducer.class); + job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); + + job.setMapOutputKeyClass(Text.class); + job.setMapOutputValueClass(IntWritable.class); + FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); - System.exit(job.waitForCompletion(true) ? 0 : 1); + + return job.waitForCompletion(true) ? 0 : 1; + } + + public static void main(String[] args) throws Exception { + final int returnCode = ToolRunner.run(new Configuration(), new greg.WordCount(), args); + System.exit(returnCode); } } diff --git a/src/main/java/greg/WordsOrder.java b/src/main/java/greg/WordsOrder.java deleted file mode 100644 index f0772d6..0000000 --- a/src/main/java/greg/WordsOrder.java +++ /dev/null @@ -1,13 +0,0 @@ -package greg; - -import org.apache.hadoop.io.IntWritable; -import org.apache.hadoop.io.LongWritable; -import org.apache.hadoop.io.Text; -import org.apache.hadoop.mapreduce.Mapper; - -public class WordsOrder { - - public static class MyMapper extends Mapper { - - } -} From 2a2e7a38ede6806c3d23cce3bb1e8b6b1d4dc846 Mon Sep 17 00:00:00 2001 From: greg Date: Fri, 9 Feb 2018 02:03:27 +0300 Subject: [PATCH 3/5] add output --- src/main/java/greg/SortWords.java | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/src/main/java/greg/SortWords.java b/src/main/java/greg/SortWords.java index 870b6d6..e63abdb 100644 --- a/src/main/java/greg/SortWords.java +++ b/src/main/java/greg/SortWords.java @@ -18,6 +18,8 @@ public class SortWords extends Configured implements Tool { + private static int counter = 0; + public static class MyMapper extends Mapper { @Override @@ -31,6 +33,12 @@ public static class MyReducer extends Reducer { @Override protected void reduce(IntWritable key, Iterable values, Context context) throws IOException, InterruptedException { + counter++; + if(7 == counter) { + System.out.println("\n\n\n\n\n\n\n\n\n\n\n" + + values.iterator().next().toString() + + "\t" + key.toString() + "\n\n\n\n\n\n\n\n\n\n"); + } super.reduce(key, values, context); } } From 3258c1e17b12ac9bf7d5dc269efb3b9c80151eba Mon Sep 17 00:00:00 2001 From: greg Date: Fri, 9 Feb 2018 03:03:15 +0300 Subject: [PATCH 4/5] finish StopWords --- src/main/java/greg/SortWords.java | 2 +- src/main/java/greg/StopWords.java | 102 ++++++++++++++++++++++++++++++ 2 files changed, 103 insertions(+), 1 deletion(-) create mode 100644 src/main/java/greg/StopWords.java diff --git a/src/main/java/greg/SortWords.java b/src/main/java/greg/SortWords.java index e63abdb..5eea5ac 100644 --- a/src/main/java/greg/SortWords.java +++ b/src/main/java/greg/SortWords.java @@ -54,7 +54,7 @@ public int compare(byte[] b1, int s1, int l1, byte[] b2, int s2, int l2) { @Override public int run(String[] args) throws Exception { final Job job = new Job(new Configuration(), "sort words"); - job.setJarByClass(WordCount.class); + job.setJarByClass(SortWords.class); job.setMapperClass(MyMapper.class); job.setReducerClass(MyReducer.class); diff --git a/src/main/java/greg/StopWords.java b/src/main/java/greg/StopWords.java new file mode 100644 index 0000000..2dbfc51 --- /dev/null +++ b/src/main/java/greg/StopWords.java @@ -0,0 +1,102 @@ +package greg; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.conf.Configured; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Counters; +import org.apache.hadoop.mapreduce.Job; +import org.apache.hadoop.mapreduce.Mapper; +import org.apache.hadoop.mapreduce.lib.input.TextInputFormat; +import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; +import org.apache.hadoop.util.Tool; +import org.apache.hadoop.util.ToolRunner; + +import java.io.File; +import java.io.FileNotFoundException; +import java.io.IOException; +import java.util.HashSet; +import java.util.Scanner; +import java.util.StringTokenizer; + +public class StopWords extends Configured implements Tool { + + private static HashSet stopWords = new HashSet(); + + public static void readStopWords(String stopWordsFilePath) { + try { + Scanner sc = new Scanner(new File(stopWordsFilePath)); + while(sc.hasNext()) { + stopWords.add(sc.nextLine()); + } + } catch (FileNotFoundException e) { + System.out.println(e.getMessage()); + } + } + + public static class MyMapper extends Mapper { + + @Override + protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { + StringTokenizer tokenizer = new StringTokenizer(value.toString()); + while (tokenizer.hasMoreTokens()) { + String word = tokenizer.nextToken(); + context.getCounter(MATCH_COUNTER.ALL).increment(1); + if (stopWords.contains(word)) { + System.out.println("Contains: " + word); + context.getCounter(MATCH_COUNTER.STOP_WORD).increment(1); + } + } + } + } + + public enum MATCH_COUNTER { + STOP_WORD, + ALL + } + + @Override + public int run(String[] args) throws Exception { + + final Job job = new Job(new Configuration(), "stop words"); + job.setJarByClass(StopWords.class); + + job.setMapperClass(MyMapper.class); + + TextInputFormat.addInputPath(job, new Path(args[0])); + TextOutputFormat.setOutputPath(job, new Path(args[1])); + + job.setInputFormatClass(TextInputFormat.class); + job.setOutputFormatClass(TextOutputFormat.class); +// job.setMapOutputKeyClass(IntWritable.class); +// job.setMapOutputValueClass(Text.class); + + job.setOutputKeyClass(Text.class); + job.setOutputValueClass(IntWritable.class); + + int exitCode = job.waitForCompletion(true) ? 0 : 1; + + Counters counter = job.getCounters(); + + System.out.println("%: " + + ((double) counter.findCounter(MATCH_COUNTER.STOP_WORD).getValue() + / counter.findCounter(MATCH_COUNTER.ALL).getValue())); + + return exitCode; + } + + public static void main(String[] args) throws Exception { + + if(args.length != 3) { + throw new IllegalArgumentException("Usage: hadoop jar .jar "); + } + + readStopWords(args[2]); + + final int returnCode = ToolRunner.run(new Configuration(), new greg.StopWords(), args); + + System.exit(returnCode); + } +} From 463d282cafcbf1f697100b05e351d8df638faaa6 Mon Sep 17 00:00:00 2001 From: greg Date: Fri, 9 Feb 2018 05:21:01 +0300 Subject: [PATCH 5/5] finish 4 tasks --- src/main/java/greg/NamesCount.java | 101 ++++++++++++++++++ src/main/java/greg/TextWithCountWritable.java | 77 +++++++++++++ 2 files changed, 178 insertions(+) create mode 100644 src/main/java/greg/NamesCount.java create mode 100644 src/main/java/greg/TextWithCountWritable.java diff --git a/src/main/java/greg/NamesCount.java b/src/main/java/greg/NamesCount.java new file mode 100644 index 0000000..7281ee4 --- /dev/null +++ b/src/main/java/greg/NamesCount.java @@ -0,0 +1,101 @@ +package greg; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.conf.Configured; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Job; +import org.apache.hadoop.mapreduce.Mapper; +import org.apache.hadoop.mapreduce.Reducer; +import org.apache.hadoop.mapreduce.lib.input.TextInputFormat; +import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; +import org.apache.hadoop.util.Tool; +import org.apache.hadoop.util.ToolRunner; + +import java.io.IOException; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +public class NamesCount extends Configured implements Tool { + + private static final Pattern namePattern = Pattern.compile("^[A-Z][a-z0-9]*$"); + + public static class MyMapper extends Mapper { + + @Override + protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { + final String[] data = value.toString().split("\t"); + + String inputString = data[1]; + int inputCount = Integer.valueOf(data[0]); + + context.write(new Text(inputString.toLowerCase()), new TextWithCountWritable(inputString, inputCount)); + } + } + + static class MyReducer extends Reducer{ + + @Override + protected void reduce(Text key, Iterable values, Context context) throws IOException, InterruptedException { + int allCount = 0; + int asNameCount= 0; + String asNameText = null; + + for (final TextWithCountWritable value : values){ + allCount += value.getCount(); + + if (asNameText == null){ + Matcher matcher = namePattern.matcher(value.getText()); + if (matcher.matches()){ + asNameText = value.getText(); + asNameCount = value.getCount(); + } + } + } + + if (asNameText == null){ + return; + } + + if (asNameCount / (double)allCount >= 0.995){ + context.write(new Text(asNameText), new IntWritable(asNameCount)); + } + } + } + + + @Override + public int run(String[] args) throws Exception { + + final Job job = new Job(new Configuration(), "count names"); + job.setJarByClass(NamesCount.class); + + job.setMapperClass(MyMapper.class); + job.setReducerClass(MyReducer.class); + + TextInputFormat.addInputPath(job, new Path(args[0])); + TextOutputFormat.setOutputPath(job, new Path(args[1])); + + job.setInputFormatClass(TextInputFormat.class); + //job.setOutputFormatClass(TextOutputFormat.class); + job.setMapOutputKeyClass(Text.class); + job.setMapOutputValueClass(TextWithCountWritable.class); + + job.setOutputKeyClass(Text.class); + job.setOutputValueClass(IntWritable.class); + + return job.waitForCompletion(true) ? 0 : 1; + } + + public static void main(String[] args) throws Exception { + if(args.length != 2) { + throw new IllegalArgumentException("Usage: hadoop jar .jar "); + } + + final int returnCode = ToolRunner.run(new Configuration(), new NamesCount(), args); + + System.exit(returnCode); + } +} diff --git a/src/main/java/greg/TextWithCountWritable.java b/src/main/java/greg/TextWithCountWritable.java new file mode 100644 index 0000000..7c0cd9b --- /dev/null +++ b/src/main/java/greg/TextWithCountWritable.java @@ -0,0 +1,77 @@ +package greg; + +// based on implementation of Georgiy0 (https://github.com/Georgiy0/2017-big-data), +// which is based on implementation of a-filippo (https://github.com/a-filippo/2017-big-data) + +import java.io.DataInput; +import java.io.DataOutput; +import java.io.IOException; + +import org.apache.hadoop.io.WritableComparable; + +public class TextWithCountWritable implements WritableComparable, Cloneable { + private String text; + private int count; + + @Override + public void write(DataOutput dataOutput) throws IOException { + dataOutput.writeInt(count); + dataOutput.writeUTF(text); + } + + @Override + public void readFields(DataInput dataInput) throws IOException { + count = dataInput.readInt(); + text = dataInput.readUTF(); + } + + public String getText() { + return text; + } + + public int getCount() { + return count; + } + + TextWithCountWritable(){} + + TextWithCountWritable(String text, int count) { + this.text = text; + this.count = count; + } + + @Override + public boolean equals(Object other) { + if (other == null){ + return false; + } + if (other == this){ + return true; + } + if (!(other instanceof TextWithCountWritable)){ + return false; + } + TextWithCountWritable otherMyClass = (TextWithCountWritable)other; + if (otherMyClass.count != count){ + return false; + } + if (!otherMyClass.text.equals(text)){ + return false; + } + return true; + } + + @Override + protected TextWithCountWritable clone() { + return new TextWithCountWritable(text, count); + } + + @Override + public int compareTo(TextWithCountWritable o) { + if (equals(o)){ + return 0; + } + int intCompare = Integer.compare(count, o.count); + return (intCompare == 0) ? Integer.compare(this.hashCode(), o.hashCode()) : -intCompare; + } +}