MapReduce中inputformat的切片、partition分区、outputformat的输出


1.1 切片
-- ① 切片的概念:从文件的逻辑上的进行大小的切分,一个切片多大,将来一个MapTask的处理的数据就多大。
-- ② 一个切片就会产生一个MapTask
-- ③ 切片时只考虑文件本身,不考虑数据的整体集。
-- ④ 切片大小和切块大小默认是一致的,这样设计目的为了避免将来切片读取数据的时候有跨机器的情况

1.2 InputFormat的体系结构

-- FileInputFormat InputFormat的子实现类,实现切片逻辑,实现了getSplits() 负责切片

-- TextInputFormat FileInputFormat的子实现类, 实现读取数据的逻辑createRecordReader() 返回一个RecordReader,在RecordReader中实现了读取数据的方式:安行读取。

-- CombineFileInputFormat FileInputFormat的子实现类,此类中也实现了一套切片逻辑 (处理:适用于小文件计算场景)

对于CombineFileInputFormat,只需在driver类中进行改进就可以了

// 如果不设置InputFormat,它默认用的是TextInputFormat.class
job.setInputFormatClass(CombineTextInputFormat.class);

//虚拟存储切片最大值设置4m
CombineTextInputFormat.setMaxInputSplitSize(job, 4194304);

2.自定义分区可以对K或V进行分区(和排序不一样,排序对象只能作为map输出类型的K值)

2.1 分为多少个区是在提交job任务时就已经确定的(通过确定reduceTask的个数),而一对KV输出到哪个区中是在map输出到环形缓冲区中(即内存中)实现的,默认的分区规则是通过KV对中的K的hashCode值除以分区的个数得到的余数来确定的,如分区个数为2,那么hashCode/2的余数为0或1。

自定义分区规则(如:对字母前缀a-q的分为0区,其余的分到1区),可通过实现

org.apache.hadoop.mapreduce.Partitioner 该类,并重写里面的getPartition()方法

注意,Partitioner的后面的数据类型和map输出的类型保持一致,因为分区的KV数据就是map输出的KV数据

public class FlowPartition extends Partitioner{}

2.2 然后在driver类中,进行配置

//设置ReduceTask数量
job.setNumReduceTasks(4);
//设置分区自定义实现类
job.setPartitionerClass(FlowPartition.class);

2.3 要查看partition的分区实现过程源码,可以在mapTask类中查看

3. 排序有两种方式:

排序在mapTask和reduceTask任务中都有(当存在分区时,reduceTask需要对来自不同mapTask的同一分区编码的KV对进行归并且按照K进行排序;而mapTask在环形缓冲区中就会在同一分区中按照K进行排序了)。

第一种:MapReduce中默认的排序是按照字典顺序,所以如果需要自定义排序规则,一要排序的对象必须实现writableComparable接口;二需要自定义一个类继承writableComparator并指定当前比较器对象为谁服务,即使用构造方法调用super传递要排序的对象的类对象;三就是在driver类中配置自定义的排序类。

一:要排序的对象必须继承writableComparable接口

public class flowBeans implements WritableComparable {

    private Integer upFlow;
    private Integer downFlow;
    private Integer totalFlow;

    //实现compare方法
    public int compareTo(flowBeans o) {
        return -this.getTotalFlow().compareTo(o.getTotalFlow());
    }

    //实现序列化
    public void write(DataOutput out) throws IOException {
        out.writeInt(upFlow);
        out.writeInt(downFlow);
        out.writeInt(totalFlow);
    }

    //实现反序列化,注意反序列化的顺序必须与序列化的顺序保持一致且序列化时如果是用writeInt,那么反序列化只能用readInt,不能用read
    public void readFields(DataInput in) throws IOException {
        upFlow = in.readInt();
        downFlow = in.readInt();
        totalFlow = in.readInt();
    }

    @Override
    public String toString() {
        return upFlow + "   " + downFlow + "    " + totalFlow ;
    }

    //因为是私有变量,外部无法直接获取,只能通过set和get方法来设置和获取私有变量
    public Integer getUpFlow() {
        return upFlow;
    }

    public void setUpFlow(Integer upFlow) {
        this.upFlow = upFlow;
    }

    public Integer getDownFlow() {
        return downFlow;
    }

    public void setDownFlow(Integer downFlow) {
        this.downFlow = downFlow;
    }

    public Integer getTotalFlow() {
        return totalFlow;
    }

    public void setTotalFlow(Integer totalFlow) {
        this.totalFlow = totalFlow;
    }    
}

二:自定义一个比较器类,类中可以重写上面的compare()方法,如果重写了,那么就会覆盖上面的compare()方法

public class flowWritable extends WritableComparator {

  //指定当前比较器对象为谁服务
    public flowWritable(){
        //调用super的前提是flowBeans必须继承writableComparable接口,点击super进去就可以看源码
        super(flowBeans.class,true);
    }

  //重写compare方法,这里可以不写,如果不写的话,那么这个比较器调用的就是上面的flowBeans对象在继承
   //writableComparable接口时所实现的compare方法。但是如果这里重写了compare()方法,那么会覆盖掉原来的compare()方法。
    @Override
    public int compare(WritableComparable a, WritableComparable b) {
        flowBeans aBean = (flowBeans) a;
        flowBeans bBean = (flowBeans) b;
        return aBean.getTotalFlow().compareTo(bBean.getTotalFlow());
    }
}

 三:在driver类中配置比较器

//配置比较器
        job.setSortComparatorClass(flowWritable.class);

第二种就是直接继承writableComparable接口就可以,不用再实现writableComparator类。用这种方式,hadoop会为该对象创建一个比较器,也就是上面的第二步。