如果一个类实现了MR的序列化的接口(Writable),这个类的对象可以作为Map和Reduce的输入和输出。本节就来使用序列化的方式重写之前的求每个部门总工资的例子。
1.程序代码
//Employee类:Employee.java
package serializable.totalsalary2;
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 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;
}
}
//TotalSalaryMapper类:TotalSalaryMapper.java
package serializable.totalsalary2;
import java.io.IOException;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
public class TotalSalaryMapper 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);
}
}
//TotalSalaryReducer类:TotalSalaryReducer.java
package serializable.totalsalary2;
import java.io.IOException;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.mapreduce.Reducer;
public class TotalSalaryReducer extends Reducer<LongWritable, Employee, LongWritable, LongWritable>{
/**
* reducer:将每个部门员工的工资求和
*/
@Override
protected void reduce(LongWritable key3, Iterable<Employee> value3, Context context)
throws IOException, InterruptedException {
long total = 0;
for(Employee e:value3) {
total += e.getSal();
}
context.write(key3,new LongWritable(total));
}
}
//Job类:TotalSalaryMain.java
package serializable.totalsalary2;
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 TotalSalaryMain {
public static void main(String[] args) throws Exception {
//创建Job
Job job = Job.getInstance(new Configuration());
//设置任务入口
job.setJarByClass(TotalSalaryMain.class);
//指定任务的Mapper,以及输出的数据类型
job.setMapperClass(TotalSalaryMapper.class);
job.setMapOutputKeyClass(LongWritable.class);
job.setMapOutputValueClass(Employee.class);
//指定任务的Reducer,以及输出的数据类型
job.setReducerClass(TotalSalaryReducer.class);
job.setOutputKeyClass(LongWritable.class);
job.setOutputValueClass(LongWritable.class);
//指定输入和输出目录:HDFS路径
FileInputFormat.setInputPaths(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
//执行任务
job.waitForCompletion(true);
}
}
2.打包并运行程序
- 将totalsalary2目录打包成totalsalary2.jar,并指定主类是TotalSalaryMain.java
- 将totalsalary2.jar上传到服务器,/root/input/totalsalary2.jar
- 准备测试数据HDFS:/input/emp.csv
- 执行程序:# hadoop jar /root/input/totalsalary2.jar /input/emp.csv /output/totalsalary2
- 查看输出目录:# hdfs dfs -ls /output/totalsalary2
- 查看结果:# hdfs dfs -cat /output/totalsalary2/part-r-00000
[root@bigdata ~]# 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 ~]# hadoop jar /root/input/totalsalary2.jar /input/emp.csv /output/totalsalary2
......
18/11/04 20:58:30 INFO mapreduce.Job: map 0% reduce 0%
18/11/04 20:58:34 INFO mapreduce.Job: map 100% reduce 0%
18/11/04 20:58:39 INFO mapreduce.Job: map 100% reduce 100%
18/11/04 20:58:40 INFO mapreduce.Job: Job job_1541335648702_0002 completed successfully
......[root@bigdata ~]# hdfs dfs -ls /output/totalsalary2
Found 2 items
-rw-r--r-- 1 root supergroup 0 2018-11-04 20:58 /output/totalsalary2/_SUCCESS
-rw-r--r-- 1 root supergroup 25 2018-11-04 20:58 /output/totalsalary2/part-r-00000[root@bigdata ~]# hdfs dfs -cat /output/totalsalary2/part-r-00000
10 8750
20 10875
30 9400