MapReduce 实现部门平均年龄排序 (多节点协同处理)
使用 MapReduce 实现部门平均年龄排序 (多节点协同处理)
需求: 找出所有部门的平均年龄,并按照平均年龄从高到低排序输出。
数据:
-
emp.csv
- 部门编号, 姓名, 年龄
- 20, zhaoyi, 30
- 30, liuer, 25
- 30, zhangsan, 31
- 20, lisi, 40
- 30, wangwu, 30
- 20, sunliu, 22
-
dept.txt
- 部门编号, 部门名称
- 10, 人事部
- 20, 技术部
- 30, 财务部
实现思路:
- 将 emp.csv 和 dept.txt 合并成一个输入文件,每一行记录包括员工的姓名、年龄和所在部门名称。
- 利用 MapReduce 框架,以部门名称为 key,将员工的年龄作为 value 进行映射,即输出 '<部门名称, 年龄>'。
- 利用 MapReduce 框架,对每个部门的年龄进行累加求和,同时记录该部门的员工数量,即输出 '<部门名称, <年龄总和, 员工数量>>'。
- 利用 MapReduce 框架,将每个部门的年龄总和和员工数量进行合并,计算出该部门的平均年龄,即输出 '<部门名称, 平均年龄>'。
- 对所有部门的平均年龄进行降序排序,输出结果。
MapReduce 实现代码:
- Mapper1
public class DeptAvgAgeMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
private final static IntWritable age = new IntWritable();
private Text dept = new Text();
public void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
String line = value.toString();
String[] fields = line.split(',');
dept.set(fields[3]); // 获取部门名称
age.set(Integer.parseInt(fields[2])); // 获取员工年龄
context.write(dept, age);
}
}
- Reducer1
public class DeptAvgAgeReducer1 extends Reducer<Text, IntWritable, Text, IntArrayWritable> {
private IntArrayWritable sum = new IntArrayWritable();
public void reduce(Text key, Iterable<IntWritable> values, Context context)
throws IOException, InterruptedException {
int ageSum = 0;
int count = 0;
for (IntWritable val : values) {
ageSum += val.get();
count++;
}
int[] arr = {ageSum, count};
sum.set(arr);
context.write(key, sum);
}
}
- Mapper2
public class DeptAvgAgeMapper2 extends Mapper<Text, IntArrayWritable, Text, DoubleWritable> {
private DoubleWritable avgAge = new DoubleWritable();
public void map(Text key, IntArrayWritable value, Context context)
throws IOException, InterruptedException {
int[] arr = Arrays.stream(value.get())
.mapToInt(i -> ((IntWritable) i).get())
.toArray();
double avg = (double) arr[0] / arr[1]; // 计算平均年龄
avgAge.set(avg);
context.write(key, avgAge);
}
}
- Reducer2
public class DeptAvgAgeReducer2 extends Reducer<Text, DoubleWritable, Text, DoubleWritable> {
private TreeMap<Double, Text> map = new TreeMap<>(Collections.reverseOrder());
public void reduce(Text key, Iterable<DoubleWritable> values, Context context)
throws IOException, InterruptedException {
double avgAge = 0;
for (DoubleWritable val : values) {
avgAge = val.get();
}
map.put(avgAge, key);
}
protected void cleanup(Context context) throws IOException, InterruptedException {
for (Map.Entry<Double, Text> entry : map.entrySet()) {
context.write(entry.getValue(), new DoubleWritable(entry.getKey()));
}
}
}
代码解释:
- Mapper1 将员工信息映射为 <部门名称, 年龄> 对。
- Reducer1 对每个部门的年龄进行累加求和,并统计员工数量,输出 <部门名称, <年龄总和, 员工数量>>。
- Mapper2 将 <部门名称, <年龄总和, 员工数量>> 映射为 <部门名称, 平均年龄>。
- Reducer2 对每个部门的平均年龄进行排序,输出最终结果。
总结:
本文介绍了使用 MapReduce 框架计算所有部门的平均年龄并按平均年龄降序排序的方法,展示了如何利用 MapReduce 框架进行多节点协同处理数据。代码示例使用 Java 语言实现,可作为参考。
原文地址: https://www.cveoy.top/t/topic/o0jE 著作权归作者所有。请勿转载和采集!