使用 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, 财务部

实现思路:

  1. 将 emp.csv 和 dept.txt 合并成一个输入文件,每一行记录包括员工的姓名、年龄和所在部门名称。
  2. 利用 MapReduce 框架,以部门名称为 key,将员工的年龄作为 value 进行映射,即输出 '<部门名称, 年龄>'。
  3. 利用 MapReduce 框架,对每个部门的年龄进行累加求和,同时记录该部门的员工数量,即输出 '<部门名称, <年龄总和, 员工数量>>'。
  4. 利用 MapReduce 框架,将每个部门的年龄总和和员工数量进行合并,计算出该部门的平均年龄,即输出 '<部门名称, 平均年龄>'。
  5. 对所有部门的平均年龄进行降序排序,输出结果。

MapReduce 实现代码:

  1. 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);
    }
}
  1. 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);
    }
}
  1. 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);
    }
}
  1. 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 语言实现,可作为参考。

MapReduce 实现部门平均年龄排序 (多节点协同处理)

原文地址: https://www.cveoy.top/t/topic/o0jE 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录