【问题标题】:Apache Hadoop is not doing it's combining & reducing work in my program that it should doApache Hadoop 没有在我的程序中组合和减少它应该做的工作
【发布时间】:2016-02-27 08:11:36
【问题描述】:

我是 Apache Hadoop 的初学者,并尝试了 Apache 的 Word Counting 程序,它运行良好。但现在我想制作自己的室外温度程序来计算每日平均值。平均计算没有像我预期的那样工作;没有对数据进行合并和平均。

更具体地说,这里是我的 sample2.txt 输入文件的一部分:

25022016 00:00:00 -10.3
25022016 00:01:00 -10.3
25022016 00:02:00 -10.3
25022016 00:03:00 -10.3
...
25022016 00:59:00 -11.2

我想要的输出应该是:

25022016 7.9

这是该日期所有温度观测值的平均值。所以我有 60 个观察值,想要一个平均值。将来,我想用相同的程序在更多的日子里处理更多的观察。 1. 列是日期(文本),2. 时间,第三列是温度。温度计算在代码中以Java的浮点数据类型完成。

现在发生的是输出是:

25022016    -10.3
25022016    -10.3
25022016    -10.3
25022016    -10.3
...
25022016    -11.2

所以每个一个观察的平均值被计算(从一个数字计算一个数字的平均值)。我想要 60 次观察的平均值(一个数字)!

所以我的输入和输出文件在上面。我的 Java 代码(我在 Windows 7 -> VirtualBox -> Ubuntu 64 位上运行)如下:


package hadoop; 

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.LongWritable;
import org.apache.hadoop.io.FloatWritable;
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;
import org.apache.hadoop.util.GenericOptionsParser;

import org.apache.commons.cli.Options;

public class ProcessUnits2 
{ 
    public static class E_EMapper extends
    Mapper<Object, Text, Text, FloatWritable>
    { 
        private FloatWritable temperature = new FloatWritable();
        private Text date = new Text();       

        public void map(Object key, Text value, 
        Context context) throws IOException, InterruptedException 
        { 
            StringTokenizer dateTimeTemperatures = new StringTokenizer(value.toString());

            while(dateTimeTemperatures.hasMoreTokens()) {
                date.set(dateTimeTemperatures.nextToken());

                while(dateTimeTemperatures.hasMoreTokens()) {
                    dateTimeTemperatures.nextToken();    
                    temperature.set(Float.parseFloat(dateTimeTemperatures.nextToken()));

                    context.write(date, temperature);
                }
            }
        } 
    } 


    public static class E_EReduce extends Reducer<Text,Text,Text,FloatWritable>
    {
        private FloatWritable result = new FloatWritable();

        public void reduce( Text key, Iterable<FloatWritable> values, Context context
        ) throws IOException, InterruptedException 
        { 
            float sumTemperatures=0, averageTemperature;
            int countTemperatures=0;

            for (FloatWritable val : values) {
                sumTemperatures += val.get();
                countTemperatures++;
            } 

            averageTemperature = sumTemperatures / countTemperatures;

            result.set(averageTemperature);
            context.write(key, result);

        } 
    }  

    public static void main(String args[])throws Exception 
    { 
        Configuration conf = new Configuration();
        String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs();

        if (otherArgs.length < 2) {
            System.err.println("Usage: wordcount <in> [<in>...] <out>");
            System.exit(2);
        }
        Job job = Job.getInstance(conf, "VuorokaudenKeskilampotila");
        job.setJarByClass(ProcessUnits2.class);

        job.setMapperClass(E_EMapper.class);
        job.setCombinerClass(E_EReduce.class);
        job.setReducerClass(E_EReduce.class);
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(FloatWritable.class);
        for (int i = 0; i < otherArgs.length - 1; ++i) {
            FileInputFormat.addInputPath(job, new Path(otherArgs[i]));
        }
        FileOutputFormat.setOutputPath(job,
        new Path(otherArgs[otherArgs.length - 1]));
        job.setNumReduceTasks(0);

        System.exit(job.waitForCompletion(true) ? 0 : 1);
    } 
} 
---------------------------------------------------

Hadoop 版本是 2.7.2 和 Ubuntu 14.04 LTS。我在独立模式下运行 hadoop(最基本的设置)。

这是我用来构建程序的命令(如果有帮助?):

rm -rf output2 
javac -Xdiags:verbose -classpath hadoop-core-1.2.1.jar:/usr/local/hadoop/share/hadoop/common/lib/commons-cli-1.2.jar -d units2 ProcessUnits2.java
jar -cvf units2.jar -C units2/ .
hadoop jar units2.jar   hadoop.ProcessUnits2 input2 output2
cat output2/part-m-00000

作为一个初学者,我很困惑,在我看来,hadoop 并没有在它的默认设置中做任何合并和减少(= 平均)工作,这应该是它的最终目的。我承认我从这里和那里(示例)中挑选了代码,因为没有任何效果,我确信这只是解决问题的一小步,但我猜不出它是什么。我可以使用例如 C++ 轻松做到这一点,根本不需要任何映射减少框架,但问题是我希望基本操作能够正常工作,因此我可以继续更复杂的示例,并在最终生产使用和真正的分布式映射组合减少.

如果能提供任何帮助,我将不胜感激。我现在陷入了困境(很多很多小时......)。如果您需要任何额外的数据来帮助找到解决方案,我会发送给他们。

【问题讨论】:

  • 感谢 Jim 提出这个更好的问题!

标签: java apache hadoop mapreduce ubuntu-14.04


【解决方案1】:

您没有正确实现减速器。应该是:

public static class E_EReduce extends Reducer<Text, FloatWritable, Text, FloatWritable>
{
    @Override
    public void reduce( Text key, Iterable<FloatWritable> values, Context context) throws IOException, InterruptedException 
    { 

永远不要忘记@Override,否则编译器不会捕捉到错误。

【讨论】:

  • 我所做的更改和行为与以前相似。我想要映射器输出:
  • 新试用版;我的 5 分钟评论时间已过:我进行了更改,并且行为与以前相似。我希望映射器输出为 25022016 -10.325022016 -17.6 ...并希望减速器在同一日期作为关键,并将所有六十个温度作为我可以迭代的值。为什么没有将具有相似键(= 25022016)的温度组合在一起,我不能在一次减少调用中迭代它们?这是关于减速器设置、全局 hadoop 配置还是其他什么?无论如何感谢您的帮助!
【解决方案2】:

现在我注意到了问题所在:

行:

job.setNumReduceTasks(0);

说没有减速器。我将其更改为job.setNumReduceTasks(1);,甚至将其完全删除,现在程序运行了。为什么它在那里? => 因为在遇到麻烦时你会尽一切可能而没有时间阅读文档。

感谢所有参与的人。我继续研究这个系统。

【讨论】:

  • 完全删除它相当于默认的job.setNumReduceTasks(1);
猜你喜欢
  • 2013-10-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-12-15
  • 2011-09-28
  • 2020-07-05
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多