辅助排序和二次排序案例(GroupingComparator)
1.需求
有如下订单数据
订单id
商品id
成交金额
0000001
Pdt_01
222.8
0000001
Pdt_05
25.8
0000002
Pdt_03
522.8
0000002
Pdt_04
122.4
0000002
Pdt_05
722.4
0000003
Pdt_01
222.8
0000003
Pdt_02
33.8
现在需要求出每一个订单中最贵的商品。
2.数据准备
GroupingComparator.txt
1Pdt_01 222.8 2 Pdt_05 722.4 3 Pdt_05 25.8 4 Pdt_01 222.8 5 Pdt_01 33.8 6 Pdt_03 522.8 7 Pdt_04 122.4
输出数据预期:

3 222.8
part-r-00000.txt

2 722.4
part-r-00001.txt

1 222.8
part-r-00002.txt
3.分析
(1)利用“订单id和成交金额”作为key,可以将map阶段读取到的所有订单数据按照id分区,按照金额排序,发送到reduce。
(2)在reduce端利用groupingcomparator将订单id相同的kv聚合成组,然后取第一个即是最大值。

4.实现
定义订单信息OrderBean
1package com.xyg.mapreduce.order; 2 3import java.io.DataInput; 4import java.io.DataOutput; 5import java.io.IOException; 6import org.apache.hadoop.io.WritableComparable; 7 8public class OrderBean implements WritableComparable<OrderBean> { 9 10 private int order_id; // 订单id号 11 private double price; // 价格 12 13 public OrderBean() { 14 super(); 15 } 16 17 public OrderBean(int order_id, double price) { 18 super(); 19 this.order_id = order_id; 20 this.price = price; 21 } 22 23 @Override 24 public void write(DataOutput out) throws IOException { 25 out.writeInt(order_id); 26 out.writeDouble(price); 27 } 28 29 @Override 30 public void readFields(DataInput in) throws IOException { 31 order_id = in.readInt(); 32 price = in.readDouble(); 33 } 34 35 @Override 36 public String toString() { 37 return order_id + "\t" + price; 38 } 39 40 public int getOrder_id() { 41 return order_id; 42 } 43 44 public void setOrder_id(int order_id) { 45 this.order_id = order_id; 46 } 47 48 public double getPrice() { 49 return price; 50 } 51 52 public void setPrice(double price) { 53 this.price = price; 54 } 55 56 // 二次排序 57 @Override 58 public int compareTo(OrderBean o) { 59 60 int result = order_id > o.getOrder_id() ? 1 : -1; 61 62 if (order_id > o.getOrder_id()) { 63 result = 1; 64 } else if (order_id < o.getOrder_id()) { 65 result = -1; 66 } else { 67 // 价格倒序排序 68 result = price > o.getPrice() ? -1 : 1; 69 } 70 71 return result; 72 } 73}
编写OrderSortMapper处理流程
1package com.xyg.mapreduce.order; 2import java.io.IOException; 3import org.apache.hadoop.io.LongWritable; 4import org.apache.hadoop.io.NullWritable; 5import org.apache.hadoop.io.Text; 6import org.apache.hadoop.mapreduce.Mapper; 7 8public class OrderMapper extends Mapper<LongWritable, Text, OrderBean, NullWritable> { 9 OrderBean k = new OrderBean(); 10 11 @Override 12 protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { 13 // 1 获取一行 14 String line = value.toString(); // 2 截取 String[] fields = line.split("\t"); // 3 封装对象 k.setOrder_id(Integer.parseInt(fields[0])); k.setPrice(Double.parseDouble(fields[2])); // 4 写出 context.write(k, NullWritable.get()); } }
编写OrderSortReducer处理流程
1package com.xyg.mapreduce.order; 2import java.io.IOException; 3import org.apache.hadoop.io.NullWritable; 4import org.apache.hadoop.mapreduce.Reducer; 5 6public class OrderReducer extends Reducer<OrderBean, NullWritable, OrderBean, NullWritable> { 7 8 @Override 9 protected void reduce(OrderBean key, Iterable<NullWritable> values, Context context) throws IOException, InterruptedException { 10 context.write(key, NullWritable.get()); 11 } 12}
编写OrderSortDriver处理流程
1package com.xyg.mapreduce.order; 2 3import java.io.IOException; 4import org.apache.hadoop.conf.Configuration; 5import org.apache.hadoop.fs.Path; 6import org.apache.hadoop.io.NullWritable; 7import org.apache.hadoop.mapreduce.Job; 8import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; 9import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; 10 11public class OrderDriver { 12 13 public static void main(String[] args) throws Exception, IOException { 14 15 // 1 获取配置信息 16 Configuration conf = new Configuration(); 17 Job job = Job.getInstance(conf); 18 19 // 2 设置jar包加载路径 20 job.setJarByClass(OrderDriver.class); 21 22 // 3 加载map/reduce类 23 job.setMapperClass(OrderMapper.class); 24 job.setReducerClass(OrderReducer.class); 25 26 // 4 设置map输出数据key和value类型 27 job.setMapOutputKeyClass(OrderBean.class); 28 job.setMapOutputValueClass(NullWritable.class); 29 30 // 5 设置最终输出数据的key和value类型 31 job.setOutputKeyClass(OrderBean.class); 32 job.setOutputValueClass(NullWritable.class); 33 34 // 6 设置输入数据和输出数据路径 35 FileInputFormat.setInputPaths(job, new Path(args[0])); 36 FileOutputFormat.setOutputPath(job, new Path(args[1])); 37 38 // 10 设置reduce端的分组 39 job.setGroupingComparatorClass(OrderGroupingComparator.class); 40 41 // 7 设置分区 42 job.setPartitionerClass(OrderPartitioner.class); 43 44 // 8 设置reduce个数 45 job.setNumReduceTasks(3); 46 47 // 9 提交 48 boolean result = job.waitForCompletion(true); 49 System.exit(result ? 0 : 1); 50 } 51} 52 53OrderSortDriver
编写OrderSortPartitioner处理流程
1package com.xyg.mapreduce.order; 2import org.apache.hadoop.io.NullWritable; 3import org.apache.hadoop.mapreduce.Partitioner; 4 5public class OrderPartitioner extends Partitioner<OrderBean, NullWritable> { 6 7 @Override 8 public int getPartition(OrderBean key, NullWritable value, int numReduceTasks) { 9 return (key.getOrder_id() & Integer.MAX_VALUE) % numReduceTasks; 10 } 11}
编写OrderSortGroupingComparator处理流程
1package com.xyg.mapreduce.order; 2import org.apache.hadoop.io.WritableComparable; 3import org.apache.hadoop.io.WritableComparator; 4 5public class OrderGroupingComparator extends WritableComparator { 6 7 protected OrderGroupingComparator() { 8 super(OrderBean.class, true); 9 } 10 11 @SuppressWarnings("rawtypes") @Override public int compare(WritableComparable a, WritableComparable b) { OrderBean aBean = (OrderBean) a; OrderBean bBean = (OrderBean) b; int result; if (aBean.getOrder_id() > bBean.getOrder_id()) { result = 1; } else if (aBean.getOrder_id() < bBean.getOrder_id()) { result = -1; } else { result = 0; } return result; } }