zoukankan      html  css  js  c++  java
  • reduce的数目到底和哪些因素有关

     
    reduce的数目到底和哪些因素有关

    1、我们知道map的数量和文件数、文件大小、块大小、以及split大小有关,而reduce的数量跟哪些因素有关呢?

     设置mapred.tasktracker.reduce.tasks.maximum的大小可以决定单个tasktracker一次性启动reduce的数目,但是不能决定总的reduce数目。

      conf.setNumReduceTasks(4);JobConf对象的这个方法可以用来设定总的reduce的数目,看下Job Counters的统计:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1. Job Counters   
    2.     Data-local map tasks=2  
    3.     Total time spent by all maps waiting after reserving slots (ms)=0  
    4.     Total time spent by all reduces waiting after reserving slots (ms)=0  
    5.     SLOTS_MILLIS_MAPS=10695  
    6.     SLOTS_MILLIS_REDUCES=29502  
    7.     Launched map tasks=2  
    8.     Launched reduce tasks=4  

     确实启动了4个reduce:看下输出:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1. diegoball@diegoball:~/IdeaProjects/test/build/classes$ hadoop fs -ls  /user/diegoball/join_ou1123  
    2. 11/03/25 15:28:45 INFO security.Groups: Group mapping impl=org.apache.hadoop.security.ShellBasedUnixGroupsMapping; cacheTimeout=300000  
    3. 11/03/25 15:28:45 WARN conf.Configuration: mapred.task.id is deprecated. Instead, use mapreduce.task.attempt.id  
    4. Found 5 items  
    5. -rw-r--r--   1 diegoball supergroup          2011-03-25 15:28 /user/diegoball/join_ou1123/_SUCCESS  
    6. -rw-r--r--   1 diegoball supergroup        124 2011-03-25 15:27 /user/diegoball/join_ou1123/part-00000  
    7. -rw-r--r--   1 diegoball supergroup          2011-03-25 15:27 /user/diegoball/join_ou1123/part-00001  
    8. -rw-r--r--   1 diegoball supergroup        214 2011-03-25 15:28 /user/diegoball/join_ou1123/part-00002  
    9. -rw-r--r--   1 diegoball supergroup          2011-03-25 15:28 /user/diegoball/join_ou1123/part-00003  

     只有2个reduce在干活。为什么呢?

    shuffle的过程,需要根据key的值决定将这条<K,V> (map的输出),送到哪一个reduce中去。送到哪一个reduce中去靠调用默认的org.apache.hadoop.mapred.lib.HashPartitioner的getPartition()方法来实现。
    HashPartitioner类:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1. package org.apache.hadoop.mapred.lib;  
    2.   
    3. import org.apache.hadoop.classification.InterfaceAudience;  
    4. import org.apache.hadoop.classification.InterfaceStability;  
    5. import org.apache.hadoop.mapred.Partitioner;  
    6. import org.apache.hadoop.mapred.JobConf;  
    7.   
    8. /** Partition keys by their {@link Object#hashCode()}.  
    9.  */  
    10. @InterfaceAudience.Public  
    11. @InterfaceStability.Stable  
    12. public class HashPartitioner<K2, V2> implements Partitioner<K2, V2> {  
    13.   
    14.   public void configure(JobConf job) {}  
    15.   
    16.   /** Use {@link Object#hashCode()} to partition. */  
    17.   public int getPartition(K2 key, V2 value,  
    18.                           int numReduceTasks) {  
    19.     return (key.hashCode() & Integer.MAX_VALUE) % numReduceTasks;  
    20.   }  
    21. }  

     numReduceTasks的值在JobConf中可以设置。默认的是1:显然太小。
       这也是为什么默认的设置中总启动一个reduce的原因。

       返回与运算的结果和numReduceTasks求余。

       Mapreduce根据这个返回结果决定将这条<K,V>,送到哪一个reduce中去。

    key传入的是LongWritable类型,看下这个LongWritable类的hashcode()方法:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1. public int hashCode() {  
    2.    return (int)value;  
    3.  }  

     简简单单的返回了原值的整型值。

     因为getPartition(K2 key, V2 value,int numReduceTask)返回的结果只有2个不同的值,所以最终只有2个reduce在干活。
     

     HashPartitioner是默认的partition类,我们也可以自定义partition类 :

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1.  package com.alipay.dw.test;  
    2.   
    3. import org.apache.hadoop.io.IntWritable;  
    4. import org.apache.hadoop.mapred.JobConf;  
    5. import org.apache.hadoop.mapred.Partitioner;  
    6.   
    7. /** 
    8.  * Created by IntelliJ IDEA. 
    9.  * User: diegoball 
    10.  * Date: 11-3-10 
    11.  * Time: 下午5:26 
    12.  * To change this template use File | Settings | File Templates. 
    13.  */  
    14. public class MyPartitioner implements Partitioner<IntWritable, IntWritable> {  
    15.     public int getPartition(IntWritable key, IntWritable value, int numPartitions) {  
    16.         /* Pretty ugly hard coded partitioning function. Don't do that in practice, it is just for the sake of understanding. */  
    17.         int nbOccurences = key.get();  
    18.         if (nbOccurences > 20051210)  
    19.             return 0;  
    20.         else  
    21.             return 1;  
    22.     }  
    23.   
    24.     public void configure(JobConf arg0) {  
    25.   
    26.     }  
    27. }  

     仅仅需要覆盖getPartition()方法就OK。通过:
    conf.setPartitionerClass(MyPartitioner.class);
    可以设置自定义的partition类。
    同样由于之返回2个不同的值0,1,不管conf.setNumReduceTasks(4);设置多少个reduce,也同样只会有2个reduce在干活。

    由于每个reduce的输出key都是经过排序的,上述自定义的Partitioner还可以达到排序结果集的目的:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1. 11/03/25 15:24:49 WARN conf.Configuration: mapred.task.id is deprecated. Instead, use mapreduce.task.attempt.id  
    2. Found 5 items  
    3. -rw-r--r--   1 diegoball supergroup          2011-03-25 15:23 /user/diegoball/opt.del/_SUCCESS  
    4. -rw-r--r--   1 diegoball supergroup      24546 2011-03-25 15:23 /user/diegoball/opt.del/part-00000  
    5. -rw-r--r--   1 diegoball supergroup      10241 2011-03-25 15:23 /user/diegoball/opt.del/part-00001  
    6. -rw-r--r--   1 diegoball supergroup          2011-03-25 15:23 /user/diegoball/opt.del/part-00002  
    7. -rw-r--r--   1 diegoball supergroup          2011-03-25 15:23 /user/diegoball/opt.del/part-00003  

     part-00000和part-00001是这2个reduce的输出,由于使用了自定义的MyPartitioner,所有key小于20051210的的<K,V>都会放到第一个reduce中处理,key大于20051210就会被放到第二个reduce中处理。
    每个reduce的输出key又是经过key排序的,所以最终的结果集降序排列。


    但是如果使用上面自定义的partition类,又conf.setNumReduceTasks(1)的话,会怎样? 看下Job Counters:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1. Job Counters   
    2.     Data-local map tasks=2  
    3.     Total time spent by all maps waiting after reserving slots (ms)=0  
    4.     Total time spent by all reduces waiting after reserving slots (ms)=0  
    5.     SLOTS_MILLIS_MAPS=16395  
    6.     SLOTS_MILLIS_REDUCES=3512  
    7.     Launched map tasks=2  
    8.     Launched reduce tasks=1  

      只启动了一个reduce。
      (1)、 当setNumReduceTasks( int a) a=1(即默认值),不管Partitioner返回不同值的个数b为多少,只启动1个reduce,这种情况下自定义的Partitioner类没有起到任何作用。
      (2)、 若a!=1:
       a、当setNumReduceTasks( int a)里 a设置小于Partitioner返回不同值的个数b的话:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1. public int getPartition(IntWritable key, IntWritable value, int numPartitions) {  
    2.     /* Pretty ugly hard coded partitioning function. Don't do that in practice, it is just for the sake of understanding. */  
    3.     int nbOccurences = key.get();  
    4.     if (nbOccurences < 20051210)  
    5.         return 0;  
    6.     if (nbOccurences >= 20051210 && nbOccurences < 20061210)  
    7.         return 1;  
    8.     if (nbOccurences >= 20061210 && nbOccurences < 20081210)  
    9.         return 2;  
    10.     else  
    11.         return 3;  
    12. }  
     

      同时设置setNumReduceTasks( 2)。

     于是抛出异常:

    [java] view plaincopy在CODE上查看代码片派生到我的代码片
     
    1.  11/03/25 17:03:41 INFO mapreduce.Job: Task Id : attempt_201103241018_0023_m_000000_1, Status : FAILED  
    2. ava.io.IOException: Illegal partition for 20110116 (3)  
    3. at org.apache.hadoop.mapred.MapTask$MapOutputBuffer.collect(MapTask.java:900)  
    4. at org.apache.hadoop.mapred.MapTask$OldOutputCollector.collect(MapTask.java:508)  
    5. at com.alipay.dw.test.KpiMapper.map(Unknown Source)  
    6. at com.alipay.dw.test.KpiMapper.map(Unknown Source)  
    7. at org.apache.hadoop.mapred.MapRunner.run(MapRunner.java:54)  
    8. at org.apache.hadoop.mapred.MapTask.runOldMapper(MapTask.java:397)  
    9. at org.apache.hadoop.mapred.MapTask.run(MapTask.java:330)  
    10. at org.apache.hadoop.mapred.Child$4.run(Child.java:217)  
    11. at java.security.AccessController.doPrivileged(Native Method)  
    12. at javax.security.auth.Subject.doAs(Subject.java:396)  
    13. at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:742)  
    14. at org.apache.hadoop.mapred.Child.main(Child.java:211)   

     某些key没有找到所对应的reduce去处。原因是只启动了a个reduce。
     
       b、当setNumReduceTasks( int a)里 a设置大于Partitioner返回不同值的个数b的话,同样会启动a个reduce,但是只有b个redurce上会得到数据。启动的其他的a-b个reduce浪费了。

       c、理想状况是a=b,这样可以合理利用资源,负载更均衡。

  • 相关阅读:
    Zabbix poller processes more than 75% busy
    标签无效 "/zabbix_export/date": "YYYY-MM-DDThh:mm:ssZ" 预计。
    Received empty response from Zabbix Agent at[172.16.1.51]. Assuming that agent dropped connection because of access permissions
    ERROR 1045 (28000): Access denied for user 'root'@'localhost' (using password: NO)
    Received empty response from Zabbix Agent at [172.16.1.7]...
    Too many open files
    Object.defineProperty之observe实现
    深拷贝 deepAssign
    提交操作自动遮蔽实现之ajax
    js检测输入域的值是否变化
  • 原文地址:https://www.cnblogs.com/bendantuohai/p/4767849.html
Copyright © 2011-2022 走看看