【问题标题】:cleaning data using mapreduce program使用 mapreduce 程序清理数据
【发布时间】:2015-03-18 17:23:10
【问题描述】:

我有一个 30 行的数据。我正在尝试使用 mapreduce 程序清理数据。数据正在正确清理,但在 30 行中只有一行显示。我猜记录阅读器不是在这里逐行阅读。您能否检查我的代码并让我知道问题出在哪里。我是 hadoop 新手。

数据:-

 1  Vlan154.DEL-ISP-COR-SWH-002.mantraonline.com (61.95.250.140)  0.460 ms  0.374 ms  0.351 ms
 2  202.56.223.213 (202.56.223.213)  39.718 ms  39.511 ms  39.559 ms
 3  202.56.223.17 (202.56.223.17)  39.714 ms  39.724 ms  39.628 ms
 4  125.21.167.153 (125.21.167.153)  41.114 ms  40.001 ms  39.457 ms
 5  203.208.190.65 (203.208.190.65)  120.340 ms  71.384 ms  71.346 ms
 6  ge-0-1-0-0.sngtp-dr1.ix.singtel.com (203.208.149.158)  71.493 ms ge-0-1-2-0.sngtp-dr1.ix.singtel.com (203.208.149.210)  71.183 ms ge-0-1-0-0.sngtp-dr1.ix.singtel.com (203.208.149.158)  71.739 ms
 7  ge-0-0-0-0.sngtp-ar3.ix.singtel.com (203.208.182.2)  80.917 ms ge-2-0-0-0.sngtp-ar3.ix.singtel.com (203.208.183.20)  71.550 ms ge-1-0-0-0.sngtp-ar3.ix.singtel.com (203.208.182.6)  71.534 ms
 8  203.208.151.26 (203.208.151.26)  141.716 ms 203.208.145.190 (203.208.145.190)  134.740 ms 203.208.151.26 (203.208.151.26)  142.453 ms
 9  219.158.3.225 (219.158.3.225)  138.774 ms  157.205 ms  157.123 ms
10  219.158.4.69 (219.158.4.69)  156.865 ms  157.044 ms  156.845 ms
11  202.96.12.62 (202.96.12.62)  157.109 ms  160.294 ms  159.805 ms
12  61.148.3.58 (61.148.3.58)  159.521 ms  178.088 ms  160.004 ms
     MPLS Label=33 CoS=5 TTL=1 S=0
13  202.106.48.18 (202.106.48.18)  199.730 ms  181.263 ms  181.300 ms
14  * * *
15  * * *
16  * * *
17  * * *
18  * * *
19  * * *
20  * * *
21  * * *
22  * * *
23  * * *

mapreduce 程序:-

公共类 TraceRouteDataCleaning {

/**
 * @param args
 * @throws IOException 
 * @throws InterruptedException 
 * @throws ClassNotFoundException 
 */
public static void main(String[] args) throws IOException, ClassNotFoundException, InterruptedException {

    Configuration conf = new Configuration();
    String userArgs[] = new GenericOptionsParser(conf, args).getRemainingArgs();
    if (userArgs.length < 2) {
        System.out.println("Usage: hadoop jar jarfilename mainclass input output");
        System.exit(1);
    }       
    Job job = new Job(conf, "cleaning trace route data");
    job.setJarByClass(TraceRouteDataCleaning.class);        
    job.setMapperClass(TraceRouteMapper.class);
    job.setReducerClass(TraceRouteReducer.class);       
    job.setMapOutputKeyClass(Text.class);
    job.setMapOutputValueClass(Text.class);
    job.setOutputKeyClass(Text.class);
    job.setOutputValueClass(Text.class);
    job.setInputFormatClass(TextInputFormat.class);
    job.setOutputFormatClass(TextOutputFormat.class);
    FileInputFormat.addInputPath(job, new Path(userArgs[0]));
    FileOutputFormat.setOutputPath(job, new Path(userArgs[1]));     
    System.exit(job.waitForCompletion(true) ? 0 : 1);
}   
public static class TraceRouteMapper extends Mapper<LongWritable, Text, Text, Text>{        
    StringBuilder emitValue = null;
    StringBuilder emitKey = null;
    Text kword = new Text();
    Text vword = new Text();

    public void map(LongWritable key, Text value, Context context) throws InterruptedException, IOException
     {
         // String[] cleanData;
         String lines = value.toString();   
         //deleting ms in RTT time data  
         lines = lines.replace(" ms", "");               
         String[] data = lines.split(" ");          
         emitValue = new StringBuilder(1024);
         emitKey = new StringBuilder(1024);

            if (data.length == 6) {                     
                emitKey.append(data[0]);
                emitValue.append(data[1]).append("\t").append(data[2]).append("\t").append(data[3]).append("\t").append(data[4]).append("\t").append(data[5]);
                kword.set(emitKey.toString());
                vword.set(emitValue.toString());                            
                context.write(kword, vword);                    
            }               
     }              
}   

public static class TraceRouteReducer extends Reducer<Text, Text, Text, Text>{
    Text vword = new Text();
    public void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException{

        context.write(key,vword);           
    }
}

}

【问题讨论】:

  • 您的拆分字符串是一个空格,但看起来您的数据有多个空格分隔多个字段。
  • @JeremyBeard:- 是的,这就是我在 map 方法中按空格分割的原因。
  • 在上面的代码中,输出只是第一行。其他线路没有来。
  • 你不是在你的 reducer 中通过键聚合,你有唯一的键吗?您可以删除 reducer 类,输出应该是 30 行,或者遍历 reducer 中每个键的值列表并输出键和值(每个键连接)
  • @Prahalad:- 我也尝试过不使用减速器类。输出只是第一行。使用减速器也没有输出变化。预期输出是前 13 行。

标签: hadoop mapreduce hdfs hadoop-streaming hadoop2


【解决方案1】:

根据您的要求,您的减速器类应该放在下面的第一件事。如果您的密钥没有发出多个文本,则选择第一个减速器,否则选择第二个。

public static class TraceRouteReducer extends Reducer<Text, Text, Text, Text>{
Text vword = new Text();
public void reduce(Text key, Text values, Context context) throws IOException, InterruptedException{

    vword=values;

    /*
 for (Iterator iterator = values.iterator(); iterator.hasNext();) {

    vword.set(iterator.next().toString());
    System.out.println("printing " +vword.toString());

}*/

    context.write(key,vword);           


}
 }

   ----------or------------

public static class TraceRouteReducer extends Reducer<Text, Text, Text, Text>{
Text vword = new Text();
  public void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException{

   for (Iterator iterator = values.iterator(); iterator.hasNext();) {

    vword.set(iterator.next().toString());
    context.write(key,vword);  


}          


}
}


second in your mapper you are splitting based on space.but not feasible as of my knowledge. split based on   "\\s+"  regular expression.

   String[] data = lines.split("\\s+");  

【讨论】:

  • Mapper 输出本身没有给出完整的结果。它只给出数据的第一行,其余数据在其中不可见。
  • 在你的映射器代码中给出 if(data.length>5)..split 基于正则表达式..我上面提到的
猜你喜欢
  • 2011-04-11
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-12-22
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多