Skip to content
Open

#1 #2

Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
.idea
/input
/output*
*/target/**
79 changes: 79 additions & 0 deletions 2017-big-data.iml
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
<?xml version="1.0" encoding="UTF-8"?>
<module org.jetbrains.idea.maven.project.MavenProjectsManager.isMavenModule="true" type="JAVA_MODULE" version="4">
<component name="NewModuleRootManager" LANGUAGE_LEVEL="JDK_1_8">
<output url="file://$MODULE_DIR$/target/classes" />
<output-test url="file://$MODULE_DIR$/target/test-classes" />
<content url="file://$MODULE_DIR$">
<sourceFolder url="file://$MODULE_DIR$/src/main/java" isTestSource="false" />
<excludeFolder url="file://$MODULE_DIR$/target" />
</content>
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-client:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-common:2.6.0" level="project" />
<orderEntry type="library" name="Maven: com.google.guava:guava:11.0.2" level="project" />
<orderEntry type="library" name="Maven: commons-cli:commons-cli:1.2" level="project" />
<orderEntry type="library" name="Maven: org.apache.commons:commons-math3:3.1.1" level="project" />
<orderEntry type="library" name="Maven: xmlenc:xmlenc:0.52" level="project" />
<orderEntry type="library" name="Maven: commons-httpclient:commons-httpclient:3.1" level="project" />
<orderEntry type="library" name="Maven: commons-codec:commons-codec:1.4" level="project" />
<orderEntry type="library" name="Maven: commons-io:commons-io:2.4" level="project" />
<orderEntry type="library" name="Maven: commons-net:commons-net:3.1" level="project" />
<orderEntry type="library" name="Maven: commons-collections:commons-collections:3.2.1" level="project" />
<orderEntry type="library" name="Maven: commons-logging:commons-logging:1.1.3" level="project" />
<orderEntry type="library" name="Maven: log4j:log4j:1.2.17" level="project" />
<orderEntry type="library" name="Maven: commons-lang:commons-lang:2.6" level="project" />
<orderEntry type="library" name="Maven: commons-configuration:commons-configuration:1.6" level="project" />
<orderEntry type="library" name="Maven: commons-digester:commons-digester:1.8" level="project" />
<orderEntry type="library" name="Maven: commons-beanutils:commons-beanutils:1.7.0" level="project" />
<orderEntry type="library" name="Maven: commons-beanutils:commons-beanutils-core:1.8.0" level="project" />
<orderEntry type="library" name="Maven: org.slf4j:slf4j-api:1.7.5" level="project" />
<orderEntry type="library" name="Maven: org.slf4j:slf4j-log4j12:1.7.5" level="project" />
<orderEntry type="library" name="Maven: org.codehaus.jackson:jackson-core-asl:1.9.13" level="project" />
<orderEntry type="library" name="Maven: org.codehaus.jackson:jackson-mapper-asl:1.9.13" level="project" />
<orderEntry type="library" name="Maven: org.apache.avro:avro:1.7.4" level="project" />
<orderEntry type="library" name="Maven: com.thoughtworks.paranamer:paranamer:2.3" level="project" />
<orderEntry type="library" name="Maven: org.xerial.snappy:snappy-java:1.0.4.1" level="project" />
<orderEntry type="library" name="Maven: com.google.protobuf:protobuf-java:2.5.0" level="project" />
<orderEntry type="library" name="Maven: com.google.code.gson:gson:2.2.4" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-auth:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.httpcomponents:httpclient:4.2.5" level="project" />
<orderEntry type="library" name="Maven: org.apache.httpcomponents:httpcore:4.2.4" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.server:apacheds-kerberos-codec:2.0.0-M15" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.server:apacheds-i18n:2.0.0-M15" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.api:api-asn1-api:1.0.0-M20" level="project" />
<orderEntry type="library" name="Maven: org.apache.directory.api:api-util:1.0.0-M20" level="project" />
<orderEntry type="library" name="Maven: org.apache.curator:curator-framework:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.curator:curator-client:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.curator:curator-recipes:2.6.0" level="project" />
<orderEntry type="library" name="Maven: com.google.code.findbugs:jsr305:1.3.9" level="project" />
<orderEntry type="library" name="Maven: org.htrace:htrace-core:3.0.4" level="project" />
<orderEntry type="library" name="Maven: org.apache.zookeeper:zookeeper:3.4.6" level="project" />
<orderEntry type="library" name="Maven: org.apache.commons:commons-compress:1.4.1" level="project" />
<orderEntry type="library" name="Maven: org.tukaani:xz:1.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-hdfs:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.mortbay.jetty:jetty-util:6.1.26" level="project" />
<orderEntry type="library" name="Maven: io.netty:netty:3.6.2.Final" level="project" />
<orderEntry type="library" name="Maven: xerces:xercesImpl:2.9.1" level="project" />
<orderEntry type="library" name="Maven: xml-apis:xml-apis:1.3.04" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-app:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-common:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-client:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-server-common:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-shuffle:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.fusesource.leveldbjni:leveldbjni-all:1.8" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-api:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-core:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-yarn-common:2.6.0" level="project" />
<orderEntry type="library" name="Maven: javax.xml.bind:jaxb-api:2.2.2" level="project" />
<orderEntry type="library" name="Maven: javax.xml.stream:stax-api:1.0-2" level="project" />
<orderEntry type="library" name="Maven: javax.activation:activation:1.1" level="project" />
<orderEntry type="library" name="Maven: javax.servlet:servlet-api:2.5" level="project" />
<orderEntry type="library" name="Maven: com.sun.jersey:jersey-core:1.9" level="project" />
<orderEntry type="library" name="Maven: com.sun.jersey:jersey-client:1.9" level="project" />
<orderEntry type="library" name="Maven: org.codehaus.jackson:jackson-jaxrs:1.9.13" level="project" />
<orderEntry type="library" name="Maven: org.codehaus.jackson:jackson-xc:1.9.13" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-mapreduce-client-jobclient:2.6.0" level="project" />
<orderEntry type="library" name="Maven: org.apache.hadoop:hadoop-annotations:2.6.0" level="project" />
</component>
</module>
11 changes: 11 additions & 0 deletions MEASUREMENTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
Выполнялись замеры на ноутбуке Xiaomi MI Air 13.3 2016 в режиме высокой производительности.

Какие имеются результаты:

WordCount job'a (без combiner):

1 час 54 минуты (грубый замер по аналоговым часам)

WordCount job'a (используя combiner):

1 час 17 минут (грубый замер по аналоговым часам)
4 changes: 2 additions & 2 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,8 @@
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>1.6</source>
<target>1.6</target>
<source>1.8</source>
<target>1.8</target>
</configuration>
</plugin>
</plugins>
Expand Down
82 changes: 82 additions & 0 deletions src/main/java/clhost/task1/WordCount.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
package clhost.task1;

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 WordCountMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
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 {
System.out.println("#map");
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 WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> {

@Override
public void reduce(final Text key, final Iterable<IntWritable> 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(WordCount.WordCountMapper.class);
job.setReducerClass(WordCount.WordCountReducer.class);
job.setCombinerClass(WordCount.WordCountReducer.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("/home/clhost/proj/2017-big-data/wordcountoutput"));

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);
}
}
85 changes: 85 additions & 0 deletions src/main/java/clhost/task2/FrequencyWordCounter.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
package clhost.task2;

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;

// подаем на вход результат работы джобы WordCount
public class FrequencyWordCounter extends Configured implements Tool {

public static class FrequencyMapper extends Mapper<LongWritable, Text, IntWritable, Text> {
private final Text word = new Text();
private final IntWritable count = new IntWritable();

@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
String[] tokens = value.toString().split("\t");

word.set(tokens[0]);
count.set((-1) * Integer.parseInt(tokens[1]));
context.write(count, word);
}
}


public static class FrequencyReducer extends Reducer<IntWritable, Text, Text, IntWritable> {
private final IntWritable count = new IntWritable();
private static final int SEEK_POSITION = 7;
private int counter = 0;

@Override
protected void reduce(IntWritable key, Iterable<Text> values, Context context) throws IOException, InterruptedException {
System.out.println("#reduce");
count.set(Integer.parseInt(key.toString()) * (-1));

for (Text value : values) {
if (counter == SEEK_POSITION - 1) {
context.write(value, count);
}
counter++;
}
}
}


@Override
public int run(final String[] args) throws Exception {
Job job = new Job(getConf(), "Frequency Word Count");
job.setJarByClass(getClass());

TextInputFormat.addInputPath(job,
new Path("/home/clhost/proj/2017-big-data/wordcountoutput/part-r-00000"));
job.setInputFormatClass(TextInputFormat.class);

job.setMapperClass(FrequencyMapper.class);
job.setReducerClass(FrequencyReducer.class);

job.setMapOutputKeyClass(IntWritable.class);
job.setMapOutputValueClass(Text.class);

TextOutputFormat.setOutputPath(job,
new Path("/home/clhost/proj/2017-big-data/frequencycountoutput"));
job.setOutputFormatClass(TextOutputFormat.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(IntWritable.class);

return job.waitForCompletion(true) ? 0 : 1;
}

public static void main(final String[] args) throws Exception {
final int returnCode = ToolRunner.run(new Configuration(), new FrequencyWordCounter(), args);
System.exit(returnCode);
}
}
112 changes: 112 additions & 0 deletions src/main/java/clhost/task3/StopWordCounter.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
package clhost.task3;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.conf.Configured;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.DoubleWritable;
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;

import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.HashSet;
import java.util.Set;
import java.util.StringTokenizer;
import java.util.stream.Collectors;
import java.util.stream.Stream;

public class StopWordCounter extends Configured implements Tool {

public static class StopWordMapper extends Mapper<LongWritable, Text, IntWritable, IntWritable> {
private HashSet<String> set = initSet();

@Override
public void map(final LongWritable key, final Text value, final Context context) throws IOException, InterruptedException {
System.out.println("#map");
final String line = value.toString();
final StringTokenizer tokenizer = new StringTokenizer(line);

int count_valid = 0;
int count_all = 0;
while (tokenizer.hasMoreTokens()) {
String word = tokenizer.nextToken();

if (set.contains(word.trim())) {
count_valid++;
}
count_all++;
}
context.write(new IntWritable(count_valid), new IntWritable(count_all));
}

private HashSet<String> initSet() {
String path = "/home/clhost/proj/2017-big-data/stop_words_en.txt";
Set<String> set = new HashSet<>();

try (Stream<String> stream = Files.lines(Paths.get(path))) {
set = stream
.collect(Collectors.toSet());
} catch (IOException e) {
e.printStackTrace();
}
return (HashSet<String>) set;
}
}


public static class StopWordReducer extends Reducer<IntWritable, IntWritable, IntWritable, DoubleWritable> {

@Override
protected void reduce(IntWritable key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException {
System.out.println("#reduce");
// приходит 2 значения
int count_valid = key.get();
int count_all = 0;

for (IntWritable ci : values) {
count_all = ci.get(); // вернет одно значение
}

// return (1, %)
context.write(new IntWritable(1), new DoubleWritable(((double) count_valid) / count_all));
}
}


@Override
public int run(final String[] args) throws Exception {
final Configuration conf = this.getConf();
final Job job = Job.getInstance(conf, "Stop Word Count");
job.setJarByClass(StopWordCounter.class);

job.setMapperClass(StopWordCounter.StopWordMapper.class);
job.setReducerClass(StopWordCounter.StopWordReducer.class);

job.setOutputKeyClass(IntWritable.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("/home/clhost/proj/2017-big-data/stopwordcountoutput"));

return job.waitForCompletion(true) ? 0 : 1;
}

public static void main(final String[] args) throws Exception {
final int returnCode = ToolRunner.run(new Configuration(), new StopWordCounter(), args);
System.exit(returnCode);
}
}
Loading