Hadoop案例(八)辅助排序和二次排序案例(GroupingComparator)

辅助排序和二次排序案例(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; } }
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang