第一章.InputFormat数据输入

$06[Hadoop-MapResource框架原理(上)] - 图1

1.切片与MapTask并行度决定机制

  1. # MapTesk并行度决定机制
  2. - 数据块:Block是HDFS物理上把数据分成一块一块
  3. - 数据切片只是在逻辑上对输入进行切片,并不会在磁盘上将其切成片进行存储

$06[Hadoop-MapResource框架原理(上)] - 图2

2.Job提交流程

$06[Hadoop-MapResource框架原理(上)] - 图3

  1. # FileInputFormat切片源码解析
  2. 1. 程序先找到你数据存储的目录
  3. 2. 开始遍历处理(规划切片)目录下的每一个文件
  4. 3. 遍历第一个文件ss.txt
  5. a. 获取文件大小 fs.sizeOf(ss.txt)
  6. b. 计算切片大小
  7. c. 默认情况下,切片大小等于块的大小
  8. d. 开始切,每次切片时,要判断切完剩下的部分是否大于块的1.1倍,不大于1.1倍就划分一块切片
  9. e. 将切片信息写到一个切片规划文件中
  10. f. 整个切片的核心过程在getSpilt()方法中完成
  11. g. InputSpilt只记录了切片的元数据信息,比如起始位置,长度以及所在的节点列表等
  12. 4. 提交切片规划文件到YARN上,Yarn上的MrAPPMaster就可以根据切片规划文件计算开启MapTask个数

3.FileInputFormat切片机制

  1. # 切片机制
  2. 1. 简单的按照文件的内容长度进行切片
  3. 2. 切片大小默认等于块的大小
  4. 3. 切片时不考虑数据集整体而是逐个针对每一个文件单独切片

$06[Hadoop-MapResource框架原理(上)] - 图4

4.CombineTextInputFormat切片机制

  1. 框架默认的TextInputFormat切片机制是对任务按文件规划切片,不管文件多小,都会是一个单独的切片,都会交给一个MapTask,这样如果有大量小文件,就会产生大量的MapTask,处理效率极其低下
  2. CombineTextInputFormat用于小文件过多的场景,它可以将多个小文件从逻辑上规划到一个切片中,这样多个小文件就可以交给一个MapTask处理
  1. # 虚拟存储切片最大值设置
  2. CombineTextInputFormat.setMaxInputSplitSize(job,4194304);//4M

切片机制

$06[Hadoop-MapResource框架原理(上)] - 图5

5.CombineTextInputFormat案例实操

  1. # 需求
  2. - 将输入的大量小文件合并成一个切片统一处理
  3. - 输入数据:准备四个小文件
  4. - 期望一个切片处理四个文件

实现过程:只需将驱动类中增加如下代码即可

  1. //设置虚拟切片的最大值
  2. CombineTextInputFormat.setMaxInputSplitSize(job,4194304);
  3. //设置使用其他的InputFormat
  4. job.setInputFormatClass(CombineTextInputFormat.class);

$06[Hadoop-MapResource框架原理(上)] - 图6

第二章.MapReducer工作流程

](https://imgtu.com/i/WEKWRO)

$06[Hadoop-MapResource框架原理(上)] - 图7

$06[Hadoop-MapResource框架原理(上)] - 图8

  1. # 上面的流程是整个MapReduce最全工作流程,但是shuffle过程只是从第7步开始到第16步结束,具体Shuffle过程详解如下:
  2. - (1)MapTask收集我们的map()方法输出的kv对,放到内存缓冲区中
  3. - (2)从内存缓冲区不断溢出本地磁盘文件,可能会溢出多个文件
  4. - (3)多个溢出文件会被合并成大的溢出文件
  5. - (4)在溢出过程及合并的过程中,都要调用Partitioner进行分区和针对key进行排序
  6. - (5)ReduceTask根据自己的分区号,去各个MapTask机器上取相应的结果分区数据
  7. - (6)ReduceTask会取到同一个分区的来自不同MapTask的结果文件,ReduceTask会将这些文件再进行合并(归并排序)
  8. - (7)合并成大文件后,Shuffle的过程也就结束了,后面进入ReduceTask的逻辑运算过程(从文件中取出一个一个的键值对Group,调用用户自定义的reduce()方法)
  1. # 注意
  2. 1. shuffle中的缓冲区大小可以通过参数调整,参数: io.sort.mb默认100M
  3. 2. shuffle缓冲区的大小会影响到MapReduce程序的执行效率,原则上说,缓冲区越大,磁盘io的次数越少,执行速度就越快

第三章.Shuffle 机制

Map方法之后,Reduce方法之前的数据处理过程称之为Shuffle

1.Shuffle机制

$06[Hadoop-MapResource框架原理(上)] - 图9

2.Partition分区

  1. # 1.问题引出
  2. 要求将统计结果按照输出条件输出到不同文件中(分区)
  3. # 2.默认Partitioner分区
  4. 1.如果有多个ReduceTask时的默认分区:
  5. public class HashPartitioner<K, V> extends Partitioner<K, V> {
  6. /*
  7. 作用 :返回分区号
  8. K : map输出的key
  9. V : map输出的value
  10. numReduceTasks : ReduceTask的数量
  11. */
  12. public int getPartition(K key, V value,
  13. int numReduceTasks) {
  14. /*
  15. 分区号 = key.hashCode() % numReduceTasks
  16. key.hashCode() & Integer.MAX_VALUE : 作用是用来保证结果为正数
  17. */
  18. return (key.hashCode() & Integer.MAX_VALUE) % numReduceTasks;
  19. }
  20. 2.如果只有1个ReduceTask时的默认分区:
  21. partitioner = new org.apache.hadoop.mapreduce.Partitioner<K,V>() {
  22. @Override
  23. public int getPartition(K key, V value, int numPartitions) {
  24. return partitions - 1; //0
  25. }
  26. };
  1. # 3.自定义Parttitioner步骤
  2. - 自定义类继承Parttitioner,重写getPartition方法
  3. - 在job驱动中,设置自定义parttitioner
  4. job.setParttitionerClass(CustomPartitioner.class)
  5. - 自定义Partition后,要根据自定义Partitioner的逻辑设置相应数量的ReduceTask
  6. job.setNumReduceTasks(5)

$06[Hadoop-MapResource框架原理(上)] - 图10

3.Partition分区案例实操

需求:将统计结果按照手机归属地不同输出到不同文件中

需求分析

$06[Hadoop-MapResource框架原理(上)] - 图11

在之前案例的基础上,新增加一个分区类

  1. package com.atguigu.mr.partition;
  2. import org.apache.hadoop.io.Text;
  3. import org.apache.hadoop.mapreduce.Partitioner;
  4. /*
  5. 自定义分区类
  6. 1.自定义类并继承Partitioner
  7. Partitioner<Key,Value>
  8. Key:mapper写出的key的类型
  9. value:mapper写出的value类型
  10. 2.重写getpartition方法
  11. */
  12. public class MyPartitioner extends Partitioner <Text,FlowBean>{
  13. /*
  14. 返回分区号
  15. */
  16. @Override
  17. public int getPartition(Text text, FlowBean flowBean, int numPartitions) {
  18. String phoneNumber = text.toString();
  19. if(phoneNumber.startsWith("136")){
  20. return 0;
  21. }else if(phoneNumber.startsWith("137")){
  22. return 1;
  23. }else if(phoneNumber.startsWith("138")){
  24. return 2;
  25. }else if(phoneNumber.startsWith("139")){
  26. return 3;
  27. }else{
  28. return 4;
  29. }
  30. }
  31. }
  1. package com.atguigu.mr.partition;
  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. //设置ReduceTask的数量
  13. job.setNumReduceTasks(5);
  14. //设置自定义分区类
  15. job.setPartitionerClass(MyPartitioner.class);
  16. //给Job对象的属性赋值
  17. job.setJarByClass(FlowDriver.class);
  18. job.setMapperClass(FlowMapper.class);
  19. job.setReducerClass(FlowReducer.class);
  20. //2.3设置Mapper输出的k,v的类型
  21. job.setMapOutputKeyClass(Text.class);
  22. job.setMapOutputValueClass(FlowBean.class);
  23. //2.4设置最终输出的K,V的类型(在这是Reducer输出的K,V类型)
  24. job.setOutputKeyClass(Text.class);
  25. job.setOutputValueClass(FlowBean.class);
  26. //2.5设置数据的输入和输出路径
  27. FileInputFormat.setInputPaths(job,new Path("C:\\io\\input4"));
  28. //注意:输出路径不能存在!!!!!!!!!!!!
  29. FileOutputFormat.setOutputPath(job,new Path("C:\\io\\output4"));
  30. //3.提交Job对象(执行Job)
  31. /*
  32. 参数 :是否打印进度
  33. 返回值 :如果是true则表示job执行成功
  34. */
  35. job.waitForCompletion(true);
  36. }
  37. }

$06[Hadoop-MapResource框架原理(上)] - 图12