第一章.shuffle机制

1.WritableComparable排序

  1. # 排序概述
  2. - 排序是MapReduce框架中最重要的操作之一
  3. - MapTask和ReduceTask均会对数据按照key进行排序该操作属于hadoop的默认行为,任何应用程序中的数据均会被排序
  4. 而不管逻辑上是否需要
  5. - 默认排序是按照字典顺序排序,且实现该排序的方法是快速排序
  6. - 对于MapTask,它会将处理的结果暂时放在环形缓冲区内,当环形缓冲区使用率达到一定阈值后,再对缓冲区中的数据进行一次快速排序,并将这些数据溢写到磁盘上,而当数据处理完毕后,它会对磁盘上的所有文件进行归并排序
  7. - 对于ReduceTask,它从每个MapTask上远程拷贝相应的数据文件,如果文件大小超过一定阈值,则溢写磁盘上,否则储存在内存中,如果磁盘上文件数目达到一定阈值,则进行一次归并排序以生成一个更大文件如果内存中文件大小或者数目超过一定阈值,则进行一次合并后将数据溢写到磁盘上,当所有的数据拷贝完毕后,ReduceTask统一对内存和磁盘上的所有数据进行一次归并排序

案例实操: 统计文件中手机号,上行流量,下行流量,总流量,并再对总流量进行排序

$07[Hadoop-MapReduce框架原理(下)] - 图1

代码实现

  1. FlowBean对象在之前的基础上增加了比较功能
  1. package com.atguigu.mr.comparable;
  2. import org.apache.hadoop.io.WritableComparable;
  3. import java.io.DataInput;
  4. import java.io.DataOutput;
  5. import java.io.IOException;
  6. /*
  7. 自定义的类的对象作为key
  8. 1.自定义的类 需要实现WritableComparable接口
  9. 2. WritableComparable继承了Writable和Comparable
  10. 3.重写Writable和Comparable接口中的抽象方法
  11. */
  12. public class FlowBean implements WritableComparable<FlowBean> {
  13. private long upFlow;
  14. private long downFlow;
  15. private long sumFlow;
  16. public FlowBean(){
  17. }
  18. public FlowBean(long upFlow, long downFlow,long sumFlow) {
  19. this.upFlow=upFlow;
  20. this.downFlow=downFlow;
  21. this.sumFlow =sumFlow;
  22. }
  23. /**
  24. * @Description: 当序列化时调用此方法
  25. * @Param: [dataOutput]
  26. * @return: void
  27. * @Author: jcsune
  28. * @Date: 2021/7/12
  29. */
  30. @Override
  31. public void write(DataOutput out) throws IOException {
  32. out.writeLong(upFlow);
  33. out.writeLong(downFlow);
  34. out.writeLong(sumFlow);
  35. }
  36. /**
  37. * @Description: 当反序列化时调用此方法,注意反序列化时读取数据的顺序要和序列化时写数据的顺序保持一致
  38. * @Param: [in]
  39. * @return: void
  40. * @Author: jcsune
  41. * @Date: 2021/7/12
  42. */
  43. @Override
  44. public void readFields(DataInput in) throws IOException {
  45. upFlow = in.readLong();
  46. downFlow = in.readLong();
  47. sumFlow = in.readLong();
  48. }
  49. /*
  50. 在该方法中指定按照哪个属性进行排序
  51. */
  52. @Override
  53. public int compareTo(FlowBean o){
  54. //按照总流量进行排序
  55. return Long.compare(this.sumFlow,o.sumFlow);
  56. }
  57. public long getUpFlow() {
  58. return upFlow;
  59. }
  60. public void setUpFlow(long upFlow) {
  61. this.upFlow = upFlow;
  62. }
  63. public long getDownFlow() {
  64. return downFlow;
  65. }
  66. public void setDownFlow(long downFlow) {
  67. this.downFlow = downFlow;
  68. }
  69. public long getSumFlow() {
  70. return sumFlow;
  71. }
  72. public void setSumFlow(long sumFlow) {
  73. this.sumFlow = sumFlow;
  74. }
  75. @Override
  76. public String toString() {
  77. return upFlow + " " + downFlow + " " +sumFlow;
  78. }
  79. }
  1. Mapper类
  1. package com.atguigu.mr.comparable;
  2. import org.apache.hadoop.io.LongWritable;
  3. import org.apache.hadoop.io.Text;
  4. import org.apache.hadoop.mapreduce.Mapper;
  5. import java.io.IOException;
  6. /*
  7. 注意:一定要让FlowBean作为key,因为只能对key进行排序
  8. */
  9. public class FlowMapper extends Mapper<LongWritable,Text,FlowBean,Text> {
  10. @Override
  11. protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
  12. //将value转成字符串
  13. String line = value.toString();
  14. //将数据切割
  15. String[] flowInfo = line.split("\t");
  16. //封装k,v
  17. FlowBean outkey =new FlowBean(Long.parseLong(flowInfo[1]),Long.parseLong(flowInfo[2]),Long.parseLong(flowInfo[3]));
  18. Text outValue = new Text(flowInfo[0]);
  19. //写出k,v
  20. context.write(outkey,outValue);
  21. }
  22. }
  1. Reducer类
  1. package com.atguigu.mr.comparable;
  2. import org.apache.hadoop.io.Text;
  3. import org.apache.hadoop.mapreduce.Reducer;
  4. import java.io.IOException;
  5. /*
  6. 只是将k,v进行调换
  7. 注意:一定要将泛型调换顺序
  8. */
  9. public class FlowReducer extends Reducer< FlowBean,Text,Text, FlowBean> {
  10. @Override
  11. protected void reduce(FlowBean key, Iterable<Text> values, Context context) throws IOException, InterruptedException {
  12. //遍历所有的value
  13. for(Text outKey: values){
  14. //将k,v写出去(原来的k变成v,原来的v变成k)
  15. context.write(outKey,key);
  16. }
  17. }
  18. }
  1. Driver类
  1. package com.atguigu.mr.comparable;
  2. import org.apache.hadoop.conf.Configuration;
  3. import org.apache.hadoop.fs.Path;
  4. import org.apache.hadoop.io.Text;
  5. import org.apache.hadoop.mapreduce.Job;
  6. import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
  7. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  8. public class FlowDriver {
  9. public static void main(String[] args) throws Exception{
  10. //创建Job对象
  11. Job job = Job.getInstance(new Configuration());
  12. //给Job对象的属性赋值
  13. job.setMapperClass(FlowMapper.class);
  14. job.setReducerClass(FlowReducer.class);
  15. //2.3设置Mapper输出的k,v的类型
  16. job.setMapOutputKeyClass(FlowBean.class);
  17. job.setMapOutputValueClass(Text.class);
  18. //2.4设置最终输出的K,V的类型(在这是Reducer输出的K,V类型)
  19. job.setOutputKeyClass(Text.class);
  20. job.setOutputValueClass(FlowBean.class);
  21. //2.5设置数据的输入和输出路径
  22. FileInputFormat.setInputPaths(job,new Path("C:\\io\\input5"));
  23. //注意:输出路径不能存在!!!!!!!!!!!!
  24. FileOutputFormat.setOutputPath(job,new Path("C:\\io\\output5"));
  25. //3.提交Job对象(执行Job)
  26. /*
  27. 参数 :是否打印进度
  28. 返回值 :如果是true则表示job执行成功
  29. */
  30. job.waitForCompletion(true);
  31. }
  32. }
  1. 执行程序,查看结果

$07[Hadoop-MapReduce框架原理(下)] - 图2

2.WritableComparable扩展(区内排序)

需求:要求每个省份手机号输出的文件中按照总流量内部排序

  1. 增加自定义分区类
  1. package com.atguigu.mr.comparable2;
  2. import org.apache.hadoop.io.Text;
  3. import org.apache.hadoop.mapreduce.Partitioner;
  4. /*
  5. 自定义分区类
  6. key:mapper输出的key
  7. value:mapper输出的value
  8. */
  9. public class MyPartitioner extends Partitioner<FlowBean, Text> {
  10. @Override
  11. public int getPartition(FlowBean flowBean, Text text, int numPartitions) {
  12. //将Text转成String
  13. String phoneNumber = text.toString();
  14. //匹配手机号
  15. if(phoneNumber.startsWith("136")){
  16. return 0;
  17. }else if(phoneNumber.startsWith("137")){
  18. return 1;
  19. }else if(phoneNumber.startsWith("137")){
  20. return 2;
  21. }else if(phoneNumber.startsWith("138")){
  22. return 3;
  23. }else {
  24. return 4;
  25. }
  26. }
  27. }
  1. 在驱动类中添加分区类
  1. package com.atguigu.mr.comparable2;
  2. import org.apache.hadoop.conf.Configuration;
  3. import org.apache.hadoop.fs.Path;
  4. import org.apache.hadoop.io.Text;
  5. import org.apache.hadoop.mapreduce.Job;
  6. import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
  7. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  8. public class FlowDriver {
  9. public static void main(String[] args) throws Exception{
  10. //创建Job对象
  11. Job job = Job.getInstance(new Configuration());
  12. //设置使用自定义分区类
  13. job.setPartitionerClass(MyPartitioner.class);
  14. //设置ReduceTask的数量
  15. job.setNumReduceTasks(5);
  16. //给Job对象的属性赋值
  17. job.setMapperClass(FlowMapper.class);
  18. job.setReducerClass(FlowReducer.class);
  19. //2.3设置Mapper输出的k,v的类型
  20. job.setMapOutputKeyClass(FlowBean.class);
  21. job.setMapOutputValueClass(Text.class);
  22. //2.4设置最终输出的K,V的类型(在这是Reducer输出的K,V类型)
  23. job.setOutputKeyClass(Text.class);
  24. job.setOutputValueClass(FlowBean.class);
  25. //2.5设置数据的输入和输出路径
  26. FileInputFormat.setInputPaths(job,new Path("C:\\io\\input6"));
  27. //注意:输出路径不能存在!!!!!!!!!!!!
  28. FileOutputFormat.setOutputPath(job,new Path("C:\\io\\output6"));
  29. //3.提交Job对象(执行Job)
  30. /*
  31. 参数 :是否打印进度
  32. 返回值 :如果是true则表示job执行成功
  33. */
  34. job.waitForCompletion(true);
  35. }
  36. }

$07[Hadoop-MapReduce框架原理(下)] - 图3

3.Combiner合并及案例实现

  1. - Combiner是MR程序中MR程序中Mapper和Reducer之外的一种组件
  2. - Combiner组件的父类就是Reducer
  3. - Combiner与Reducer的区别在于运行的位置
  4. - Combiner是在每一个MapTask所在的节点运行
  5. - Reducer是接收全局所有Mapper的输出结果
  6. - Combiner的意义就是对每一个MapTask的输出进行局部汇总,以减少网络传输量
  7. - Combiner能够应用的前提下是不能影响最终的业务逻辑

案例实操: 统计过程中对每一个MapTask的输出进行局部汇总,以减少网络传输量即采用Combiner功能

未使用Combiner

$07[Hadoop-MapReduce框架原理(下)] - 图4

使用Combiner

  1. package com.atguigu.mr.combiner;
  2. import org.apache.hadoop.conf.Configuration;
  3. import org.apache.hadoop.fs.Path;
  4. import org.apache.hadoop.io.IntWritable;
  5. import org.apache.hadoop.io.Text;
  6. import org.apache.hadoop.mapreduce.Job;
  7. import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
  8. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  9. public class WCDriver {
  10. public static void main(String[] args) throws Exception{
  11. //1.创建job对象
  12. Configuration conf = new Configuration();//可以在该对象中做一些设置
  13. Job job = Job.getInstance();
  14. //设置Combiner
  15. job.setCombinerClass(WCCombiner.class);
  16. //2.给job对象的属性赋值
  17. //设置jar加载路径(如果是本地模式,可以不设置)
  18. job.setJarByClass(WCDriver.class);
  19. //设置Mapper和Reduce类
  20. job.setMapperClass(WCMapper.class);
  21. job.setReducerClass(WCReducer.class);
  22. //设置Mapper输出的k,v类型
  23. job.setMapOutputKeyClass(Text.class);
  24. job.setMapOutputValueClass(IntWritable.class);
  25. //设置最终输出的K,V类型
  26. job.setOutputKeyClass(Text.class);
  27. job.setOutputValueClass(IntWritable.class);
  28. //设置数据的输入和输出路径
  29. FileInputFormat.setInputPaths(job,new Path("C:\\io\\input7"));
  30. //注意输出路径不能存在
  31. FileOutputFormat.setOutputPath(job,new Path("C:\\io\\output7"));
  32. //3.提交job对象
  33. /*
  34. 参数:是否打印进度
  35. 返回值:如果是true则表示job执行成功
  36. */
  37. boolean b = job.waitForCompletion(true);
  38. System.exit(b ? 0 : 1);
  39. }
  40. }

$07[Hadoop-MapReduce框架原理(下)] - 图5

4.outputFormat数据输出

  1. # outputFormat接口实现类
  2. - outputFormat是MapReduce输出的基类,所有实现MapReduce输出都实现了OutputFormat接口
  3. - 文本输出TextOutputFormat:默认的输出格式是TextOutputFormat,它把每条记录写为文本行,它的键和值可以是任意类型,因为TextOutputFormat调用toString()方法把他们转换为字符串
  4. - SequenceFileOutputFormat:将SequerceFileOutputFormat输出作为后续MapReduce任务的输入,这便是好的一种输出格式,因为它的格式紧凑,很容易被压缩
  5. - 自定义OutputFormat
  6. 1. 自定义一个类继承FileOutputFormat
  7. 2. 改写RecordWriter,具体改写输出数据的方法write()

案例实操

  1. 需求:过滤输入的log日志,包含atguigu的网站输出到atguigu.txt,不包含atguigu的网站输出到other.txt
  2. 需求分析

$07[Hadoop-MapReduce框架原理(下)] - 图6

MyRecordWriter.java

  1. package com.atguigu.mr.output;
  2. import org.apache.hadoop.fs.FSDataOutputStream;
  3. import org.apache.hadoop.fs.FileSystem;
  4. import org.apache.hadoop.fs.Path;
  5. import org.apache.hadoop.io.IOUtils;
  6. import org.apache.hadoop.io.LongWritable;
  7. import org.apache.hadoop.io.Text;
  8. import org.apache.hadoop.mapreduce.RecordWriter;
  9. import org.apache.hadoop.mapreduce.TaskAttemptContext;
  10. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  11. import java.io.IOException;
  12. /*
  13. 自定义RecordWriter类
  14. */
  15. public class MyRecordWriter extends RecordWriter<LongWritable, Text> {
  16. private FSDataOutputStream atguigu;
  17. private FSDataOutputStream other;
  18. /*
  19. 创建流
  20. */
  21. public MyRecordWriter(TaskAttemptContext job) {
  22. try{
  23. //1.创建流
  24. FileSystem fs = FileSystem.get(job.getConfiguration());
  25. //2.创建流--输出流
  26. //获取输出路径
  27. Path outputPath = FileOutputFormat.getOutputPath(job);
  28. atguigu =fs.create(new Path(outputPath,"atguigu.txt"));
  29. other = fs.create(new Path(outputPath,"other.txt"));
  30. }catch (Exception e){
  31. //终止程序运行
  32. throw new RuntimeException(e.getMessage());
  33. }
  34. }
  35. /*
  36. 写数据
  37. */
  38. @Override
  39. public void write(LongWritable key, Text value) throws IOException, InterruptedException {
  40. //先判断网址有没有包含atguigu
  41. String address = value.toString() + "\n";
  42. if(address.contains("atguigu")){
  43. //写到atguigu.txt
  44. atguigu.write(address.getBytes());
  45. }else{
  46. //写到other.txt
  47. other.write(address.getBytes());
  48. }
  49. }
  50. /*
  51. 在该方法中用来关闭资源---当写操作结束以后才会调用此方法(框架调用的)
  52. */
  53. @Override
  54. public void close(TaskAttemptContext context) throws IOException, InterruptedException {
  55. IOUtils.closeStream(atguigu);
  56. IOUtils.closeStream(other);
  57. }
  58. }

MyOutputFormat.java

  1. package com.atguigu.mr.output;
  2. import org.apache.hadoop.io.LongWritable;
  3. import org.apache.hadoop.io.Text;
  4. import org.apache.hadoop.mapreduce.RecordWriter;
  5. import org.apache.hadoop.mapreduce.TaskAttemptContext;
  6. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  7. import java.io.IOException;
  8. public class MyOutputFormat extends FileOutputFormat<LongWritable,Text> {
  9. /*
  10. 自定义OutputFormat类
  11. 1. 自定义一个类并继承FileOutputFormat,因为继承此类只需要重写getRecordWriter方法而此方法正是我们所需要的
  12. 2. FileOutputFormat泛型的类型Reducer输出的K,V类型
  13. 没有Reducer是Mapper输出的K,V
  14. 没有Mapper也没有那么是Inputformat输出的K,V类型
  15. */
  16. @Override
  17. public RecordWriter<LongWritable, Text> getRecordWriter(TaskAttemptContext job) throws IOException, InterruptedException {
  18. return new MyRecordWriter(job);
  19. }
  20. }

OutputDriver.java

  1. package com.atguigu.mr.output;
  2. import org.apache.hadoop.conf.Configuration;
  3. import org.apache.hadoop.fs.Path;
  4. import org.apache.hadoop.io.LongWritable;
  5. import org.apache.hadoop.io.Text;
  6. import org.apache.hadoop.mapreduce.Job;
  7. import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
  8. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  9. public class OutputDriver {
  10. public static void main(String[] args) throws Exception {
  11. Job job = Job.getInstance(new Configuration());
  12. job.setOutputKeyClass(LongWritable.class);
  13. job.setOutputValueClass(Text.class);
  14. //设置自定义的outputFormat类--如果不指定使用TextOutputFormat
  15. job.setOutputFormatClass(MyOutputFormat.class);
  16. FileInputFormat.setInputPaths(job,new Path("C:\\io\\input8"));
  17. FileOutputFormat.setOutputPath(job,new Path("C:\\io\\output8"));
  18. job.waitForCompletion(true);
  19. }
  20. }

输出结果

$07[Hadoop-MapReduce框架原理(下)] - 图7

5.Join多种应用

1. ReduceJoin工作原理

  1. - Map端的主要工作:为来自不同表或文件的key/value对,打标签以区别不同的来源记录,然后用连接字段作为key,其余部分和新加的标志作为value,最后进行输出
  2. - Reduce端的主要工作:在Reduce端以连接字段作为key的分组已经完成.我们只需要在每一个分组当中将那些来源于不同文件的记录(在Map阶段已经打标志分开,)最后进行合并就OK了

2. 案例实操

$07[Hadoop-MapReduce框架原理(下)] - 图8

$07[Hadoop-MapReduce框架原理(下)] - 图9

MyCompartor.java

  1. package com.atguigu.mr.reducejoin;
  2. import org.apache.hadoop.io.WritableComparable;
  3. import org.apache.hadoop.io.WritableComparator;
  4. /*
  5. 自定义分组
  6. 1.自定义一个类并继承WritableComparator
  7. 2.调用父类的有参构造器
  8. 3.重写compare方法,并在该方法中指定分组的方式
  9. */
  10. public class MyCompartor extends WritableComparator {
  11. /*
  12. WritableComparator(Class<? extends WritableComparable> keyClass,
  13. boolean createInstances)
  14. keyClass : Key对象所在的类的运行时类的对象
  15. createInstances : 是否实例化
  16. */
  17. public MyCompartor(){
  18. super(OrderBean.class,true);
  19. }
  20. /*
  21. 重写compare方法,并在该方法中指定分组的方式
  22. 作用:在这里按照pid 进行分组
  23. */
  24. @Override
  25. public int compare(WritableComparable a, WritableComparable b) {
  26. //向下转型
  27. OrderBean o1 = (OrderBean) a;
  28. OrderBean o2 = (OrderBean) b;
  29. return o1.getPid().compareTo(o2.getPid());
  30. }
  31. }

OrderBean.java

  1. package com.atguigu.mr.reducejoin;
  2. import org.apache.hadoop.io.WritableComparable;
  3. import java.io.DataInput;
  4. import java.io.DataOutput;
  5. import java.io.IOException;
  6. public class OrderBean implements WritableComparable<OrderBean> {
  7. private String id;
  8. private String pid;
  9. private String amount;
  10. private String pname;
  11. public OrderBean(){
  12. }
  13. public OrderBean(String id,String pid,String amount,String pname){
  14. this.id = id;
  15. this.pid = pid;
  16. this.amount = amount;
  17. this.pname = pname;
  18. }
  19. public String getId() {
  20. return id;
  21. }
  22. public void setId(String id) {
  23. this.id = id;
  24. }
  25. public String getPid() {
  26. return pid;
  27. }
  28. public void setPid(String pid) {
  29. this.pid = pid;
  30. }
  31. public String getAmount() {
  32. return amount;
  33. }
  34. public void setAmount(String amount) {
  35. this.amount = amount;
  36. }
  37. public String getPname() {
  38. return pname;
  39. }
  40. public void setPname(String pname) {
  41. this.pname = pname;
  42. }
  43. /**
  44. * @Description: 排序,先按照pid进行排序,pid相同再按照pname记性排序
  45. * @Param: [o]
  46. * @return: int
  47. * @Author: jcsune
  48. * @Date: 2021/7/15
  49. */
  50. @Override
  51. public int compareTo(OrderBean o) {
  52. int pidValue = this.pid.compareTo(o.pid);
  53. if(pidValue==0){//说明pid相同,在按照pname排序
  54. return -this.pname.compareTo(o.pname);
  55. }
  56. return pidValue;
  57. }
  58. /**
  59. * @Description: 序列化
  60. * @Param: [dataOutput]
  61. * @return: void
  62. * @Author: jcsune
  63. * @Date: 2021/7/15
  64. */
  65. @Override
  66. public void write(DataOutput dataOutput) throws IOException {
  67. dataOutput.writeUTF(id);
  68. dataOutput.writeUTF(pid);
  69. dataOutput.writeUTF(pname);
  70. dataOutput.writeUTF(amount);
  71. }
  72. /**
  73. * @Description: 反序列化
  74. * @Param: [dataInput]
  75. * @return: void
  76. * @Author: jcsune
  77. * @Date: 2021/7/15
  78. */
  79. @Override
  80. public void readFields(DataInput dataInput) throws IOException {
  81. id = dataInput.readUTF();
  82. pid = dataInput.readUTF();
  83. pname = dataInput.readUTF();
  84. amount = dataInput.readUTF();
  85. }
  86. @Override
  87. public String toString() {
  88. return id + " " + pid + " " + amount + " " + pname;
  89. }
  90. }

RJmapper.java

  1. package com.atguigu.mr.reducejoin;
  2. import org.apache.hadoop.io.LongWritable;
  3. import org.apache.hadoop.io.NullWritable;
  4. import org.apache.hadoop.io.Text;
  5. import org.apache.hadoop.mapreduce.Mapper;
  6. import org.apache.hadoop.mapreduce.lib.input.FileSplit;
  7. import java.io.IOException;
  8. public class RJMapper extends Mapper<LongWritable, Text,OrderBean, NullWritable> {
  9. private String fileName;
  10. /*
  11. 该方法在任务开始的时候只执行一次,在map方法调用前执行
  12. 初始化
  13. */
  14. @Override
  15. protected void setup(Context context) throws IOException, InterruptedException {
  16. //获取切片信息
  17. FileSplit fs = (FileSplit) context.getInputSplit();
  18. //通过切片获取文件名
  19. fileName = fs.getPath().getName();
  20. }
  21. @Override
  22. protected void cleanup(Context context) throws IOException, InterruptedException {
  23. super.cleanup(context);
  24. }
  25. @Override
  26. protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
  27. String line = value.toString();
  28. //切割数据
  29. String[] info = line.split("\t");
  30. //创建对象
  31. OrderBean bean = new OrderBean();
  32. //封装
  33. if("order.txt".equals(fileName)){
  34. bean.setId(info[0]);
  35. bean.setPid(info[1]);
  36. bean.setAmount(info[2]);
  37. //注意:一定要给Pname一个空串,因为排序的时候先按照pid排序再按照pname排序
  38. bean.setPname("");
  39. }
  40. if("pd.txt".equals(fileName)){
  41. bean.setPid(info[0]);
  42. bean.setPname(info[1]);
  43. bean.setId("");
  44. bean.setAmount("");
  45. }
  46. //写出K,V
  47. context.write(bean,NullWritable.get());
  48. }
  49. }

RJReducer.java

  1. package com.atguigu.mr.reducejoin;
  2. import org.apache.hadoop.io.NullWritable;
  3. import org.apache.hadoop.mapreduce.Reducer;
  4. import java.io.IOException;
  5. import java.util.Iterator;
  6. public class RJReducer extends Reducer<OrderBean, NullWritable,OrderBean,NullWritable> {
  7. @Override
  8. protected void reduce(OrderBean key, Iterable<NullWritable> values, Context context) throws IOException, InterruptedException {
  9. //获取迭代器对象
  10. Iterator<NullWritable> iterator = values.iterator();
  11. //获取该数据中的第一行
  12. iterator.next();
  13. //获取pname
  14. String pname = key.getPname();
  15. //将剩余的所有数据的pname进行替换
  16. while(iterator.hasNext()){
  17. iterator.next();
  18. //替换pname
  19. key.setPname(pname);
  20. //将k,v写出去
  21. context.write(key,NullWritable.get());
  22. }
  23. }
  24. }

RJDriver.java

  1. package com.atguigu.mr.reducejoin;
  2. import org.apache.hadoop.fs.Path;
  3. import org.apache.hadoop.io.NullWritable;
  4. import org.apache.hadoop.mapreduce.Job;
  5. import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
  6. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  7. public class RJDriver {
  8. public static void main(String[] args) throws Exception {
  9. Job job = Job.getInstance();
  10. //设置自定义分组类
  11. job.setGroupingComparatorClass(MyCompartor.class);
  12. job.setMapperClass(RJMapper.class);
  13. job.setReducerClass(RJReducer.class);
  14. job.setMapOutputKeyClass(OrderBean.class);
  15. job.setMapOutputValueClass(NullWritable.class);
  16. job.setOutputKeyClass(OrderBean.class);
  17. job.setOutputValueClass(NullWritable.class);
  18. FileInputFormat.setInputPaths(job,new Path("C:\\io\\input9"));
  19. FileOutputFormat.setOutputPath(job,new Path("C:\\io\\output9"));
  20. job.waitForCompletion(true);
  21. }
  22. }

3.MapJoin工作原理

  1. # 使用场景
  2. - Map join使用于一张表十分小,一张表很大的场景
  3. - 优点: 在Map端缓存多张表,提前处理业务逻辑,这样增加Map端业务,减少Reduce端数据的压力,尽可能的减少数据倾斜
  4. - 具体方法: 采用DistributedCache
  5. 1. 在Mapper的Setup阶段,将文件读取到缓存集合中
  6. 2. 在驱动函数中加载缓存
  7. // 缓存普通文件到Task运行节点
  8. job.addCacheFile(New URL("file://e:/cache/pd.txt");

4. 案例实操

需求:将商品信息表根据商品pid合并到订单数据表中

$07[Hadoop-MapReduce框架原理(下)] - 图10

OrderBean.java

  1. package com.atguigu.mr.mapjoin;
  2. import org.apache.hadoop.io.WritableComparable;
  3. import java.io.DataInput;
  4. import java.io.DataOutput;
  5. import java.io.IOException;
  6. public class OrderBean implements WritableComparable<OrderBean> {
  7. private String id;
  8. private String pid;
  9. private String amount;
  10. private String pname;
  11. public OrderBean(){
  12. }
  13. public OrderBean(String id, String pid, String amount, String pname){
  14. this.id = id;
  15. this.pid = pid;
  16. this.amount = amount;
  17. this.pname = pname;
  18. }
  19. public String getId() {
  20. return id;
  21. }
  22. public void setId(String id) {
  23. this.id = id;
  24. }
  25. public String getPid() {
  26. return pid;
  27. }
  28. public void setPid(String pid) {
  29. this.pid = pid;
  30. }
  31. public String getAmount() {
  32. return amount;
  33. }
  34. public void setAmount(String amount) {
  35. this.amount = amount;
  36. }
  37. public String getPname() {
  38. return pname;
  39. }
  40. public void setPname(String pname) {
  41. this.pname = pname;
  42. }
  43. /**
  44. * @Description: 排序,先按照pid进行排序,pid相同再按照pname记性排序
  45. * @Param: [o]
  46. * @return: int
  47. * @Author: jcsune
  48. * @Date: 2021/7/15
  49. */
  50. @Override
  51. public int compareTo(OrderBean o) {
  52. int pidValue = this.pid.compareTo(o.pid);
  53. if(pidValue==0){//说明pid相同,在按照pname排序
  54. return -this.pname.compareTo(o.pname);
  55. }
  56. return pidValue;
  57. }
  58. /**
  59. * @Description: 序列化
  60. * @Param: [dataOutput]
  61. * @return: void
  62. * @Author: jcsune
  63. * @Date: 2021/7/15
  64. */
  65. @Override
  66. public void write(DataOutput dataOutput) throws IOException {
  67. dataOutput.writeUTF(id);
  68. dataOutput.writeUTF(pid);
  69. dataOutput.writeUTF(pname);
  70. dataOutput.writeUTF(amount);
  71. }
  72. /**
  73. * @Description: 反序列化
  74. * @Param: [dataInput]
  75. * @return: void
  76. * @Author: jcsune
  77. * @Date: 2021/7/15
  78. */
  79. @Override
  80. public void readFields(DataInput dataInput) throws IOException {
  81. id = dataInput.readUTF();
  82. pid = dataInput.readUTF();
  83. pname = dataInput.readUTF();
  84. amount = dataInput.readUTF();
  85. }
  86. @Override
  87. public String toString() {
  88. return id + " " + pid + " " + amount + " " + pname;
  89. }
  90. }

MJMapper.java

  1. package com.atguigu.mr.mapjoin;
  2. import org.apache.hadoop.fs.FSDataInputStream;
  3. import org.apache.hadoop.fs.FileSystem;
  4. import org.apache.hadoop.fs.Path;
  5. import org.apache.hadoop.io.LongWritable;
  6. import org.apache.hadoop.io.NullWritable;
  7. import org.apache.hadoop.io.Text;
  8. import org.apache.hadoop.mapreduce.Mapper;
  9. import java.io.BufferedReader;
  10. import java.io.IOException;
  11. import java.io.InputStreamReader;
  12. import java.net.URI;
  13. import java.util.HashMap;
  14. import java.util.Map;
  15. public class MJMapper extends Mapper<LongWritable, Text,OrderBean, NullWritable> {
  16. private Map<String,String> map = new HashMap<String,String>();
  17. /*
  18. 在任务开始的时候,只执行一次在map方法之前执行
  19. 作用 将pd.txt中的内容缓存起来
  20. */
  21. @Override
  22. protected void setup(Context context) throws IOException, InterruptedException {
  23. FileSystem fs = null;
  24. BufferedReader br =null;
  25. try {
  26. //创建流-读取pd.txt
  27. fs= FileSystem.get(context.getConfiguration());
  28. //获取缓存文件
  29. URI[] cacheFiles = context.getCacheFiles();
  30. //创建流
  31. FSDataInputStream fis = fs.open(new Path(cacheFiles[0]));
  32. //读数据-- 一行一行的读
  33. //字符缓冲流----将字节流转成字符流再套用字符缓冲流
  34. br = new BufferedReader(new InputStreamReader(fis,"utf-8"));
  35. //读数据
  36. String line = "";
  37. while((line = br.readLine())!= null){
  38. //将数据切割
  39. String[] split = line.split("\t");
  40. //将数据存放到map中
  41. map.put(split[0],split[1]);
  42. }
  43. } catch (IOException e) {
  44. e.printStackTrace();
  45. }finally {
  46. if (br != null){
  47. br.close();
  48. }
  49. if(fs != null){
  50. fs.close();
  51. }
  52. }
  53. }
  54. @Override
  55. protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
  56. String line = value.toString();
  57. //切割数据
  58. String[] info = line.split("\t");
  59. //封装对象
  60. OrderBean orderBean = new OrderBean(info[0],info[1],info[2],map.get(info[1]));
  61. //写出k,v
  62. context.write(orderBean,NullWritable.get());
  63. }
  64. }

MJDriver.class

  1. package com.atguigu.mr.mapjoin;
  2. import org.apache.hadoop.conf.Configuration;
  3. import org.apache.hadoop.fs.Path;
  4. import org.apache.hadoop.io.NullWritable;
  5. import org.apache.hadoop.mapreduce.Job;
  6. import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
  7. import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
  8. import java.net.URI;
  9. public class MJDriver {
  10. public static void main(String[] args) throws Exception {
  11. Job job = Job.getInstance(new Configuration());
  12. job.addCacheFile(new URI("file:///C:/io/input10/pd.txt"));//添加缓存文件的路径,一定要注意斜杠
  13. job.setMapperClass(MJMapper.class);
  14. job.setNumReduceTasks(0);//表示没有ReducerTask
  15. job.setMapOutputKeyClass(OrderBean.class);
  16. job.setMapOutputValueClass(NullWritable.class);
  17. job.setMapOutputKeyClass(OrderBean.class);
  18. job.setMapOutputValueClass(NullWritable.class);
  19. FileInputFormat.setInputPaths(job,new Path("C:\\io\\input10\\order.txt"));
  20. FileOutputFormat.setOutputPath(job,new Path("C:\\io\\output10"));
  21. job.waitForCompletion(true);
  22. }
  23. }