Hadoop从入门到精通31:MapReduce高级功能之分区

1.什么是分区?

在进行MapReduce计算时,有时候需要把最终的输出数据分到不同的文件中,比如按照省份划分的话,需要把同一省份的数据放到一个文件中;按照性别划分的话,需要把同一性别的数据放到一个文件中等等。我们知道最终的输出数据是来自于Reducer任务。那么,如果要得到多个文件,意味着有同样数量的Reducer任务在运行。Reducer任务的数据来自于Mapper任务,也就说Mapper任务要划分数据,对于不同的数据分配给不同的Reducer任务运行。Mapper任务划分数据的过程就称作Partition。负责实现划分数据的类称作Partitioner。

分区的英文单词叫做Partition,简写为part。可以从MR任务的输出文件part-r-00000的前缀part看到这一点。MR默认情况下只有一个分区,即只有一个输出文件part-r-00000。如果设置了多个分区,那么就会在一个目录下输出多个文件:part-r-00000,part-r-00001,part-r-00002,等等。

对比日志信息:

(1)没有分区的情况(一个分区)

18/11/05 22:12:38 INFO mapreduce.Job: map 0% reduce 0%
18/11/05 22:12:41 INFO mapreduce.Job: map 100% reduce 0%
18/11/05 22:12:45 INFO mapreduce.Job: map 100% reduce 100%

(2)有分区的情况(3个分区)

18/11/05 22:12:38 INFO mapreduce.Job: map 0% reduce 0%
18/11/05 22:12:41 INFO mapreduce.Job: map 100% reduce 0%
18/11/05 22:12:45 INFO mapreduce.Job: map 100% reduce 33%
18/11/05 22:12:49 INFO mapreduce.Job: map 100% reduce 67%
18/11/05 22:12:54 INFO mapreduce.Job: map 100% reduce 100%

2.开发带有分区的MR程序

(1)开发带有分区的MR程序需要注意以下几点:

  1. 在Mapper和Reducer之间加上一个Partitioner阶段;
  2. Partitioner的输入就是Mapper的输出;
  3. 分区类需要继承自Partitioner父类;
  4. 分区类需要重载getPartition()方法;

示例:按照员工的部门号对员工数据进行分类存放(不同部门的员工输出到不同的文件中)。

//员工类:Employee.java
package demo.part;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
import org.apache.hadoop.io.Writable;
public class Employee implements Writable {
  private int empno;
  private String ename;
  private String job;
  private int mgr;
  private String hiredate;
  private int sal;
  private int comm;
  private int deptno;
  @Override
  public String toString() {
    return "["+this.ename+"\t"+this.deptno+"\t"+this.sal+"]";
  }
  @Override
  public void readFields(DataInput input) throws IOException {
    // 反序列化
    this.empno = input.readInt();
    this.ename = input.readUTF();
    this.job = input.readUTF();
    this.mgr = input.readInt();
    this.hiredate = input.readUTF();
    this.sal = input.readInt();
    this.comm = input.readInt();
    this.deptno = input.readInt();
  }
  @Override
  public void write(DataOutput output) throws IOException {
    // 序列化
    output.writeInt(this.empno);
    output.writeUTF(this.ename);
    output.writeUTF(this.job);
    output.writeInt(this.mgr);
    output.writeUTF(this.hiredate);
    output.writeInt(this.sal);
    output.writeInt(this.comm);
    output.writeInt(this.deptno);
  }
  public int getEmpno() {
    return empno;
  }
  public void setEmpno(int empno) {
    this.empno = empno;
  }
  public String getEname() {
    return ename;
  }
  public void setEname(String ename) {
    this.ename = ename;
  }
  public String getJob() {
    return job;
  }
  public void setJob(String job) {
    this.job = job;
  }
  public int getMgr() {
    return mgr;
  }
  public void setMgr(int mgr) {
    this.mgr = mgr;
  }
  public String getHiredate() {
    return hiredate;
  }
  public void setHiredate(String hiredate) {
    this.hiredate = hiredate;
  }
  public int getSal() {
    return sal;
  }
  public void setSal(int sal) {
    this.sal = sal;
  }
  public int getComm() {
    return comm;
  }
  public void setComm(int comm) {
    this.comm = comm;
  }
  public int getDeptno() {
    return deptno;
  }
  public void setDeptno(int deptno) {
    this.deptno = deptno;
  }
}
//Mapper类:EmployeePartitionMapper.java
package demo.part;
import java.io.IOException;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
public class EmployeePartitionMapper extends Mapper<LongWritable, Text, LongWritable, Employee>{
  /*
  * mapper:将读入的员工数据存到一个员工对象中
  */
  @Override
  protected void map(LongWritable key1, Text value1, Context context)
    throws IOException, InterruptedException {
    //读入一行数据:7654,MARTIN,SALESMAN,7698,1981/9/28,1250,1400,30
    String line = value1.toString();
    //分词操作
    String[] words = line.split(",");
    //创建一个员工对象
    Employee e = new Employee();
    //设置员工属性
    e.setEmpno(Integer.parseInt(words[0]));
    e.setEname(words[1]);
    e.setJob(words[2]);
    try {
      e.setMgr(Integer.parseInt(words[3]));
    }catch(Exception ex) {
      e.setMgr(0);
    }
    e.setHiredate(words[4]);
    e.setSal(Integer.parseInt(words[5]));
    try {
      e.setComm(Integer.parseInt(words[6]));
    }catch(Exception ex) {
      e.setComm(0);
    }
    e.setDeptno(Integer.parseInt(words[7]));
    //map输出
    context.write(new LongWritable(e.getDeptno()), e);
  }
}
//Partitioner类:EmployeePartitioner.java
package demo.part;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.mapreduce.Partitioner;
public class EmployeePartitioner extends Partitioner<LongWritable, Employee>{
  @Override
  public int getPartition(LongWritable k2, Employee v2, int numPart) {
    // 参数:(k2,v2)就是Mapper的输出,numPart是分区数
    // 根据该员工的部门号返回该员工所属的分区号
    int deptno = v2.getDeptno();
    if(deptno == 10){
      return 1%numPart;
    }else if(deptno ==20){
      return 2%numPart;
    }else{
      return 3%numPart;
    }
  }
}
//Reducer类:EmployeePartitionReducer.java
package demo.part;
import java.io.IOException;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.mapreduce.Reducer;
public class EmployeePartitionReducer extends Reducer<LongWritable, Employee, LongWritable, Employee>{
  @Override
  protected void reduce(LongWritable key3, Iterable<Employee> value3, Context context)
      throws IOException, InterruptedException {
    // 将分区后的结果直接输出到HDFS
    for(Employee e:value3) {
      context.write(key3,e);
    }
  }
}
//Job类:EmployeePartitionMain.java
package demo.part;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
public class EmployeePartitionMain {
  public static void main(String[] args) throws Exception {
    //创建Job
    Job job = Job.getInstance(new Configuration());
    //设置任务入口
    job.setJarByClass(EmployeePartitionMain.class);
    //指定任务的Mapper,以及输出的数据类型
    job.setMapperClass(EmployeePartitionMapper.class);
    job.setMapOutputKeyClass(LongWritable.class);
    job.setMapOutputValueClass(Employee.class);
    //指定任务的分区规则
    job.setPartitionerClass(EmployeePartitioner.class);
    //指定分区的个数
    job.setNumReduceTasks(3); 
    //指定任务的Reducer,以及输出的数据类型
    job.setReducerClass(EmployeePartitionReducer.class);
    job.setOutputKeyClass(LongWritable.class);
    job.setOutputValueClass(Employee.class);
    //指定输入和输出目录:HDFS路径
    FileInputFormat.setInputPaths(job, new Path(args[0]));
    FileOutputFormat.setOutputPath(job, new Path(args[1]));
    //执行任务
    job.waitForCompletion(true);
  }
}

(2)打包并运行程序

  1. 将demo.part目录打包成EmployeePartition.jar,并指定主类是EmployeePartitionMain.java
  2. 将EmployeePartition.jar上传到服务器,/root/input/EmployeePartition.jar
  3. 准备测试数据HDFS:/input/emp.csv
  4. 执行程序:# hadoop jar /root/input/EmployeePartition.jar /input/emp.csv /output/employeepartition
  5. 查看输出目录:# hdfs dfs -ls /output/employeepartition
  6. 查看结果:# hdfs dfs -cat /output/employeepartition/part-r-00000

[root@bigdata input]# hdfs dfs -cat /input/emp.csv
7369,SMITH,CLERK,7902,1980/12/17,800,,20
7499,ALLEN,SALESMAN,7698,1981/2/20,1600,300,30
7521,WARD,SALESMAN,7698,1981/2/22,1250,500,30
7566,JONES,MANAGER,7839,1981/4/2,2975,,20
7654,MARTIN,SALESMAN,7698,1981/9/28,1250,1400,30
7698,BLAKE,MANAGER,7839,1981/5/1,2850,,30
7782,CLARK,MANAGER,7839,1981/6/9,2450,,10
7788,SCOTT,ANALYST,7566,1987/4/19,3000,,20
7839,KING,PRESIDENT,,1981/11/17,5000,,10
7844,TURNER,SALESMAN,7698,1981/9/8,1500,0,30
7876,ADAMS,CLERK,7788,1987/5/23,1100,,20
7900,JAMES,CLERK,7698,1981/12/3,950,,30
7902,FORD,ANALYST,7566,1981/12/3,3000,,20
7934,MILLER,CLERK,7782,1982/1/23,1300,,10

[root@bigdata input]# hadoop jar EmployeePartition.jar /input/emp.csv /output/employeepartition
......
18/11/05 23:41:09 INFO mapreduce.Job: map 0% reduce 0%
18/11/05 23:41:14 INFO mapreduce.Job: map 100% reduce 0%
18/11/05 23:41:20 INFO mapreduce.Job: map 100% reduce 33%
18/11/05 23:41:22 INFO mapreduce.Job: map 100% reduce 100%
18/11/05 23:41:23 INFO mapreduce.Job: Job job_1541425571272_0004 completed successfully
......

[root@bigdata input]# hdfs dfs -ls /output/employeepartition
Found 4 items
-rw-r--r-- 1 root supergroup 0 2018-11-05 23:41 /output/employeepartition/_SUCCESS
-rw-r--r-- 1 root supergroup 114 2018-11-05 23:41 /output/employeepartition/part-r-00000
-rw-r--r-- 1 root supergroup 57 2018-11-05 23:41 /output/employeepartition/part-r-00001
-rw-r--r-- 1 root supergroup 93 2018-11-05 23:41 /output/employeepartition/part-r-00002

[root@bigdata input]# hdfs dfs -cat /output/employeepartition/part-r-00000
30 [MARTIN 30 1250]
30 [JAMES 30 950]
30 [BLAKE 30 2850]
30 [WARD 30 1250]
30 [TURNER 30 1500]
30 [ALLEN 30 1600]
[root@bigdata input]# hdfs dfs -cat /output/employeepartition/part-r-00001
10 [MILLER 10 1300]
10 [KING 10 5000]
10 [CLARK 10 2450]
[root@bigdata input]# hdfs dfs -cat /output/employeepartition/part-r-00002
20 [SCOTT 20 3000]
20 [JONES 20 2975]
20 [ADAMS 20 1100]
20 [FORD 20 3000]
20 [SMITH 20 800]

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容