diff --git a/differences between task with Combainer and without it b/differences between task with Combainer and without it new file mode 100644 index 0000000..f9c52ff --- /dev/null +++ b/differences between task with Combainer and without it @@ -0,0 +1,2 @@ +word count без комбайнера 1 час 40 мин +word count c комбайнера 1 час 10 мин diff --git a/pom.xml b/pom.xml index 16acaa5..4f6bd4a 100644 --- a/pom.xml +++ b/pom.xml @@ -1,47 +1,76 @@ - - 4.0.0 - pritykovskaya - 2017-big-data - jar - 1.0-SNAPSHOT - 2017-big-data - http://maven.apache.org + + + 4.0.0 - - UTF-8 - 2.6.0 - + MapReduce_wordCounter + MapReduce_wordCounter + 1.0-SNAPSHOT - - - org.apache.hadoop - hadoop-client - ${hadoop.version} - - + + + + UTF-8 + 2.6.0 + + + + + org.apache.hadoop + hadoop-client + ${hadoop.version} + + + + + + + maven-assembly-plugin + + + package + + single + + + + + + + true + mapReduceTask.core.task_1.WordCountJob + + + + jar-with-dependencies + + + + + org.apache.maven.plugins + maven-compiler-plugin + + 1.7 + 1.7 + + + + - - - - org.apache.maven.plugins - maven-jar-plugin - - - - pritykovskaya.WordCount - - - - - - org.apache.maven.plugins - maven-compiler-plugin - - 1.6 - 1.6 - - - - diff --git a/src/main/java/pritykovskaya/WordCount.java b/src/main/java/pritykovskaya/WordCount.java deleted file mode 100644 index b047511..0000000 --- a/src/main/java/pritykovskaya/WordCount.java +++ /dev/null @@ -1,80 +0,0 @@ -package pritykovskaya; - - -import java.io.IOException; -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.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.FileInputFormat; -import org.apache.hadoop.mapreduce.lib.input.TextInputFormat; -import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; -import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; -import org.apache.hadoop.util.Tool; -import org.apache.hadoop.util.ToolRunner; - - -public class WordCount extends Configured implements Tool { - - public static class MyMapper extends Mapper { - private static final IntWritable ONE = new IntWritable(1); - private final transient Text word = new Text(); - - @Override public void map(final LongWritable key, final Text value, final Context context) - throws IOException, InterruptedException { - final String line = value.toString(); - final StringTokenizer tokenizer = new StringTokenizer(line); - while (tokenizer.hasMoreTokens()) { - word.set(tokenizer.nextToken()); - context.write(word, ONE); - } - } - } - - - public static class MyReducer extends Reducer { - - @Override - public void reduce(final Text key, final Iterable values, final Context context) - throws IOException, InterruptedException { - int sum = 0; - for (final IntWritable val : values) { - sum += val.get(); - } - context.write(key, new IntWritable(sum)); - } - } - - - @Override public int run(final String[] args) throws Exception { - final Configuration conf = this.getConf(); - final Job job = Job.getInstance(conf, "Word Count"); - job.setJarByClass(WordCount.class); - - job.setMapperClass(MyMapper.class); - job.setReducerClass(MyReducer.class); - - job.setOutputKeyClass(Text.class); - job.setOutputValueClass(IntWritable.class); - - job.setInputFormatClass(TextInputFormat.class); - job.setOutputFormatClass(TextOutputFormat.class); - - FileInputFormat.addInputPath(job, new Path(args[0])); - FileOutputFormat.setOutputPath(job, new Path(args[1])); - - return job.waitForCompletion(true) ? 0 : 1; - } - - public static void main(final String[] args) throws Exception { - final int returnCode = ToolRunner.run(new Configuration(), new WordCount(), args); - System.exit(returnCode); - } -} diff --git a/src/main/java/task_1/WordCountJob.java b/src/main/java/task_1/WordCountJob.java new file mode 100755 index 0000000..f469176 --- /dev/null +++ b/src/main/java/task_1/WordCountJob.java @@ -0,0 +1,50 @@ +package task_1; + +import org.apache.hadoop.conf.Configured; +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.lib.input.TextInputFormat; +import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; +import org.apache.hadoop.util.Tool; +import org.apache.hadoop.util.ToolRunner; + +/** + * Created by iters on 10/21/17. + */ +public class WordCountJob extends Configured implements Tool { + public int run(String[] strings) throws Exception { + Job job = new Job(getConf(), "WordCount"); + job.setJarByClass(getClass()); + + // TextInputFormat.addInputPath(job, new Path(strings[1])); + TextInputFormat.addInputPath(job, + new Path(strings[0])); + + job.setInputFormatClass(TextInputFormat.class); + + job.setMapperClass(WordCounterMapper.class); + job.setReducerClass(WordCounterReducer.class); + job.setCombinerClass(WordCounterReducer.class); + + // TextOutputFormat.setOutputPath(job, new Path(strings[2])); + TextOutputFormat.setOutputPath(job, + new Path(strings[1])); + job.setOutputFormatClass(TextOutputFormat.class); + job.setOutputKeyClass(Text.class); + job.setOutputValueClass(IntWritable.class); + + return job.waitForCompletion(true) ? 0 : 1; + } + + public static void main(String[] args) { + try { + int exitCode = ToolRunner.run(new WordCountJob(), args); + System.exit(exitCode); + } catch (Exception e) { + System.out.println("error on start ;("); + e.printStackTrace(); + } + } +} \ No newline at end of file diff --git a/src/main/java/task_1/WordCounterMapper.java b/src/main/java/task_1/WordCounterMapper.java new file mode 100755 index 0000000..0f5fbfd --- /dev/null +++ b/src/main/java/task_1/WordCounterMapper.java @@ -0,0 +1,30 @@ +package task_1; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Mapper; + +import java.io.IOException; +import java.util.StringTokenizer; + +/** + * Created by iters on 10/22/17. + */ +public class WordCounterMapper + extends Mapper { + + private final static IntWritable one = new IntWritable(1); + private final Text word = new Text(); + + @Override + protected void map(LongWritable key, Text value, Context context) + throws IOException, InterruptedException { + StringTokenizer st = new StringTokenizer(value.toString()); + + while (st.hasMoreTokens()) { + word.set(st.nextToken()); + context.write(word, one); + } + } +} \ No newline at end of file diff --git a/src/main/java/task_1/WordCounterReducer.java b/src/main/java/task_1/WordCounterReducer.java new file mode 100755 index 0000000..98d2005 --- /dev/null +++ b/src/main/java/task_1/WordCounterReducer.java @@ -0,0 +1,25 @@ +package task_1; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Reducer; + +import java.io.IOException; + +/** + * Created by iters on 10/22/17. + */ +public class WordCounterReducer + extends Reducer { + + @Override + protected void reduce(Text key, Iterable values, Context context) + throws IOException, InterruptedException { + int sum = 0; + for (IntWritable num : values) { + sum += num.get(); + } + + context.write(key, new IntWritable(sum)); + } +} \ No newline at end of file diff --git a/src/main/java/task_2/DescendingSorting.java b/src/main/java/task_2/DescendingSorting.java new file mode 100755 index 0000000..409c1f6 --- /dev/null +++ b/src/main/java/task_2/DescendingSorting.java @@ -0,0 +1,50 @@ +package task_2; + +import org.apache.hadoop.conf.Configured; +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.lib.input.TextInputFormat; +import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; +import org.apache.hadoop.util.Tool; +import org.apache.hadoop.util.ToolRunner; + +/** + * Created by iters on 10/21/17. + */ + +public class DescendingSorting extends Configured implements Tool { + public int run(String[] strings) throws Exception { + Job job = new Job(getConf(), "DescendingSortingWords"); + job.setJarByClass(getClass()); + + TextInputFormat.addInputPath(job, + new Path(strings[0])); + job.setInputFormatClass(TextInputFormat.class); + + job.setMapperClass(WordCounterMapper.class); + job.setReducerClass(WordCounterReducer.class); + + job.setMapOutputKeyClass(IntWritable.class); + job.setMapOutputValueClass(Text.class); + + TextOutputFormat.setOutputPath(job, + new Path(strings[1])); + job.setOutputFormatClass(TextOutputFormat.class); + job.setOutputKeyClass(Text.class); + job.setOutputValueClass(IntWritable.class); + + return job.waitForCompletion(true) ? 0 : 1; + } + + public static void main(String[] args) { + try { + int exitCode = ToolRunner.run(new DescendingSorting(), args); + System.exit(exitCode); + } catch (Exception e) { + System.out.println("error on start ;("); + e.printStackTrace(); + } + } +} diff --git a/src/main/java/task_2/WordCounterMapper.java b/src/main/java/task_2/WordCounterMapper.java new file mode 100755 index 0000000..ab29755 --- /dev/null +++ b/src/main/java/task_2/WordCounterMapper.java @@ -0,0 +1,28 @@ +package task_2; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Mapper; +import java.io.IOException; +import java.util.Arrays; + +/** + * Created by iters on 10/22/17. + */ +public class WordCounterMapper + extends Mapper { + + private final Text word = new Text(); + private final IntWritable num = new IntWritable(); + + @Override + protected void map(LongWritable key, Text value, Context context) + throws IOException, InterruptedException { + String[] parts = value.toString().split("\t"); + + word.set(parts[0]); + num.set(-1 * Integer.parseInt(parts[1])); + context.write(num, word); + } +} \ No newline at end of file diff --git a/src/main/java/task_2/WordCounterReducer.java b/src/main/java/task_2/WordCounterReducer.java new file mode 100755 index 0000000..6005a70 --- /dev/null +++ b/src/main/java/task_2/WordCounterReducer.java @@ -0,0 +1,32 @@ +package task_2; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Reducer; +import java.io.IOException; + +/** + * Created by iters on 10/22/17. + */ +public class WordCounterReducer + extends Reducer { + + private final IntWritable positiveInt = new IntWritable(); + private static Text text = new Text(); + private static final int OUT_POS = 7; + private static int pos = 0; + + @Override + protected void reduce(IntWritable key, Iterable values, Context context) + throws IOException, InterruptedException { + + positiveInt.set(Integer.parseInt(key.toString()) * (-1)); + for (Text value : values) { + if (++pos == OUT_POS) { + text.set("7-ое по популярности слово:\n" + + value.toString()); + context.write(text, positiveInt); + } + } + } +} \ No newline at end of file diff --git a/src/main/java/task_3/StopWordProportion.java b/src/main/java/task_3/StopWordProportion.java new file mode 100755 index 0000000..0628ef4 --- /dev/null +++ b/src/main/java/task_3/StopWordProportion.java @@ -0,0 +1,61 @@ +package task_3; + +import org.apache.hadoop.conf.Configured; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Counters; +import org.apache.hadoop.mapreduce.Job; +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; + +/** + * Created by iters on 10/21/17. + */ +public class StopWordProportion extends Configured implements Tool { + + enum stopWordProportion { + STOP_WORD, + COMMON_WORD + } + + public int run(String[] strings) throws Exception { + Job job = new Job(getConf(), "StopWordProportion"); + job.setJarByClass(getClass()); + + TextInputFormat.addInputPath(job, new Path(strings[0])); + job.setInputFormatClass(TextInputFormat.class); + + job.setMapperClass(WordCounterMapper.class); + // job.setReducerClass(task_1.WordCounterReducer.class); + job.setNumReduceTasks(0); + + TextOutputFormat.setOutputPath(job, new Path(strings[1])); + job.setOutputFormatClass(TextOutputFormat.class); + job.setOutputKeyClass(IntWritable.class); + job.setOutputValueClass(Text.class); + + int exitCode = job.waitForCompletion(true) ? 0 : 1; + + Counters counters = job.getCounters(); + long stopWords = counters.findCounter(stopWordProportion.STOP_WORD).getValue(); + long commonWords = counters.findCounter(stopWordProportion.COMMON_WORD).getValue(); + double res = Double.parseDouble(stopWords + "") / + (commonWords + stopWords) * 100d; + + System.out.println(res + "% стоп-слов в en_articles"); + return exitCode; + } + + public static void main(String[] args) { + try { + int exitCode = ToolRunner.run(new StopWordProportion(), args); + System.exit(exitCode); + } catch (Exception e) { + System.out.println("error on start ;("); + e.printStackTrace(); + } + } +} diff --git a/src/main/java/task_3/WordCounterMapper.java b/src/main/java/task_3/WordCounterMapper.java new file mode 100755 index 0000000..a12942b --- /dev/null +++ b/src/main/java/task_3/WordCounterMapper.java @@ -0,0 +1,49 @@ +package task_3; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Mapper; + +import java.io.*; +import java.util.HashSet; +import java.util.Set; +import java.util.StringTokenizer; + +/** + * Created by iters on 10/22/17. + */ +public class WordCounterMapper + extends Mapper { + + private static Set stopWords = new HashSet<>(); + + static { + InputStream is = WordCounterMapper.class.getClassLoader().getResourceAsStream("StopWords"); + try (BufferedReader bf = new BufferedReader(new InputStreamReader(is))) { + String line; + + while ((line = bf.readLine()) != null) { + stopWords.add(line); + } + } catch (IOException e) { + e.printStackTrace(); + } + } + + @Override + protected void map(LongWritable key, Text value, Context context) + throws IOException, InterruptedException { + + StringTokenizer st = new StringTokenizer(value.toString()); + while(st.hasMoreTokens()) { + if (stopWords.contains(st.nextToken())) { + context.getCounter( + StopWordProportion.stopWordProportion.STOP_WORD).increment(1); + } else { + context.getCounter( + StopWordProportion.stopWordProportion.COMMON_WORD).increment(1); + } + } + } +} \ No newline at end of file diff --git a/src/main/java/task_4/NameCounterJob.java b/src/main/java/task_4/NameCounterJob.java new file mode 100644 index 0000000..0b104f2 --- /dev/null +++ b/src/main/java/task_4/NameCounterJob.java @@ -0,0 +1,48 @@ +package task_4; + +import org.apache.hadoop.conf.Configured; +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.lib.input.TextInputFormat; +import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; +import org.apache.hadoop.util.Tool; +import org.apache.hadoop.util.ToolRunner; + +/** + * Created by iters on 11/13/17. + */ +public class NameCounterJob extends Configured implements Tool { + + @Override + public int run(String[] strings) throws Exception { + Job job = new Job(getConf(), "NameCount"); + job.setJarByClass(getClass()); + + TextInputFormat.addInputPath(job, + new Path(strings[0])); + + job.setInputFormatClass(TextInputFormat.class); + job.setMapperClass(WordCounterMapper.class); + job.setReducerClass(NameCounterReducer.class); + + TextOutputFormat.setOutputPath(job, + new Path(strings[1])); + job.setOutputFormatClass(TextOutputFormat.class); + job.setOutputKeyClass(Text.class); + job.setOutputValueClass(IntWritable.class); + + return job.waitForCompletion(true) ? 0 : 1; + } + + public static void main(String[] args) { + try { + int exitCode = ToolRunner.run(new NameCounterJob(), args); + System.exit(exitCode); + } catch (Exception e) { + System.out.println("error on start ;("); + e.printStackTrace(); + } + } +} diff --git a/src/main/java/task_4/NameCounterReducer.java b/src/main/java/task_4/NameCounterReducer.java new file mode 100644 index 0000000..f711590 --- /dev/null +++ b/src/main/java/task_4/NameCounterReducer.java @@ -0,0 +1,42 @@ +package task_4; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Reducer; + +import java.io.IOException; + +/** + * Created by iters on 11/13/17. + */ +public class NameCounterReducer + extends Reducer { + + private final IntWritable digit = new IntWritable(); + private final Text tName = new Text(); + + @Override + protected void reduce(Text key, Iterable values, Context context) throws IOException, InterruptedException { + Double sumPositive = 0d; + Double sumNegative = 0d; + + for (IntWritable num : values) { + int res = num.get(); + + if (res == 1) { + sumPositive++; + } else { + sumNegative++; + } + } + + if (Double.compare(sumNegative / (sumPositive + sumNegative) * 100 + , 0.5) < 0) { + digit.set(sumPositive.intValue()); + tName.set(key.toString().substring(0, 1).toUpperCase() + + key.toString().substring(1, key.toString().length())); + + context.write(tName, digit); + } + } +} diff --git a/src/main/java/task_4/WordCounterMapper.java b/src/main/java/task_4/WordCounterMapper.java new file mode 100755 index 0000000..7715939 --- /dev/null +++ b/src/main/java/task_4/WordCounterMapper.java @@ -0,0 +1,42 @@ +package task_4; + +import org.apache.hadoop.io.IntWritable; +import org.apache.hadoop.io.LongWritable; +import org.apache.hadoop.io.Text; +import org.apache.hadoop.mapreduce.Mapper; + +import java.io.IOException; +import java.util.StringTokenizer; + +/** + * Created by iters on 10/22/17. + */ +public class WordCounterMapper + extends Mapper { + + private final Text word = new Text(); + + @Override + protected void map(LongWritable key, Text value, Context context) + throws IOException, InterruptedException { + StringTokenizer st = new StringTokenizer(value.toString()); + String str, temp; + + while (st.hasMoreTokens()) { + str = st.nextToken(); + temp = str; + str = str.toLowerCase(); + str = str.substring(0, 1).toUpperCase() + + str.substring(1, str.length()); + + word.set(str.toLowerCase()); + + boolean isDigit = Character.isDigit(str.charAt(0)); + if (temp.equals(str) && !isDigit) { + context.write(word, new IntWritable(1)); + } else { + context.write(word, new IntWritable(-1)); + } + } + } +} \ No newline at end of file diff --git a/src/main/resources/StopWords b/src/main/resources/StopWords new file mode 100644 index 0000000..9db7df0 --- /dev/null +++ b/src/main/resources/StopWords @@ -0,0 +1,2 @@ +hello +world \ No newline at end of file