欢迎来到安国润的个人网站

MapReduce概述

返回列表

MapReduce概述

一、MapReduce 整体定位与核心作用

1.1 MapReduce 简单介绍

MapReduce 是谷歌(Google) 在 2004 年正式提出的分布式大数据计算模型,出自工程师 Jeff Dean 与 Sanjay Ghemawat 发布的论文《MapReduce: Simplified Data Processing on Large Clusters》。

设计初衷:解决搜索引擎中大规模网页数据的并行处理

后续开源社区基于论文思想开发出 Hadoop MapReduce,成为大数据离线计算的基础,Hive 数仓的 HQL 底层运行时,都会自动翻译成 MapReduce 任务执行。

MapReduce:2004 谷歌,大数据分布式计算框架,用于海量离线数据统计;

Transformer:2017 谷歌,AI 深度学习网络结构,用于大语言模型、图像 AI。

1.2 核心思想:分而治之

MapReduce 是 Hadoop 第一代分布式计算框架,面向批处理的分布式计算框架,专门用于海量离线数据的分布式统计、清洗、聚合,是 Hive数仓工具底层真正的执行引擎。

特点:计算跟着数据走,良好的拓展性,高容错,状态监控,适合海量数据的离线批处理,降低了分布式编程的门槛

适用场景:

数据统计,如网页的pv 浏览量,uv 用户数统计

搜索引擎构建索引

海量数据查询

复杂数据分析算法实现

不适用场景:

OLAP 要求毫秒及秒级返回结果

流式计算 流式计算的输入数据是动态 的,而mapreduce的静态的

DAG 计算 多个作业有依赖关系,后一个的输入是前一个的输出,构成有向无环图

每个mapreduce作业的输出结果都会落盘,构成大量磁盘IO,导致性能非常低。一般spark做DAG

Map阶段:负责将工作任务分解为若干个子任务来并行处理,这些子任务相互独立,可以单独被执行。

Reduce阶段:负责将Map过程处理完的子任务结果合并,从而得到工作任务的最终结果。

img

MapReduce就是“任务的分解与结果的汇总”。即使用户不懂分布式计算框架的内部运行机制,但是只要能用

Map和Reduce思想描述清楚要处理的问题,就能轻松地在Hadoop集群上实现分布式计算功能。

1.3 编程模型

MapReduce 是一种编程模型,用于处理大规模数据集的并行运算。使用MapReduce执行计算任务的时候,每个任务的执行过程都会被分为两个阶段,分别是Map和Reduce,其中Map阶段用于对原始数据进行处理,Reduce阶段用于对Map阶段的结果进行汇总,得到最终结果。

在这里插入图片描述

1.4 经典实例——词频统计

1782374938159

词频统计 文本出现的次数,

第一步 split 文本切分为多个数据块,如果文本是在hdfs中,hdfs 自动完成split切分,如果不在hdfs上,会按128M一块切分,拆分多块

第二步 每一个数据块会起一个maptask,进行运算,maptask会处理成 KV的形式,得到部分结果

第三步 汇总到最终结果,把k相同的数据分发到同一个reduce节点进行聚合, 之后每个reduce中k相同的v相加在一起,得到词频统计的结果,最终输出。可以手动指定reduce个数

最难的是中间中间走shuffer阶段, k相同的如何分发到同一个reduce中,通过hash取模,k进行数字编码,对数字进行hash取模,如果有4个reduce,k就对4hash取模,实际上是除以4 取余数 0123,对应4个reduce。

假设有两个文本文件test1.txt和文件test2.txt。

使用MapReduce程序统计文件test1.txt和test2.txt中每个单词出现的次数,实现词频统计的流程。

  • 首先,MapReduce通过默认组件TextInputFormat将待处理的数据文件(如text1.txt和text2.txt),把每一行的数据都转变为键值对。其中键(Key)是指每行数据的起始偏移量,也就是每行数据开头的字符所在的位置,值(Value)是指文本文件中的每行数据。
  • 其次,调用Map()方法,将单词进行切割并进行计数,输出键值对作为Reduce阶段的输入键值对。
  • 最后,调用Reduce()方法将单词汇总、排序后,通过TextOutputFormat组件输出到结果文件中。
  • 在这里插入图片描述

二、五大阶段分步详解(每一步功能、完成工作)

整套流程分为五大阶段:Input 输入 → Map 阶段 → Shuffle 洗牌阶段 → Reduce 阶段 → Output 输出

1782375004240

  1. 文件上传到hdfs上自动完成切分为多个块
  2. 每个块会启动一个map任务,map在运行过程中,会处理为kv结果数据,
  3. shuffle map要把k相同的聚合到同一个reduce节点,对k进行reduce个数哈希取模,得到三个组文件,都数据分完组之后
  4. reduce 就开始拉取所需要的文件,进行文件合并,然后再reduce处理,一般是k相同的v进行相关的处理,词频统计就是累加。最终每个reduce输入一个结果文件,这些结果文件会存放到同一个结果目录下,这些文件加起来就是我们最后的结果

阶段 1:Input 数据输入分片阶段

  1. 完成工作
  2. 读取 HDFS 上的原始大文件(CSV、文本);
  3. 将超大文件切分成多个固定大小的数据分片(Split),每一片分配给一台机器的 MapTask 并行处理;
  4. 读取每行原始数据,封装成 <偏移量key, 单行文本value> 传递给 Map 函数。
  5. 对应实训业务:读取采集好的新发地农产品 CSV 原始数据,拆分多行数据分给多个 Map 并行解析。

阶段 2:Map 映射阶段(map () 方法)

  1. 核心功能:逐行处理原始明细数据,提取、过滤、转换字段,输出中间键值对 <K2,V2>

  2. 具体完成的工作:

    1)读取分片内每一行原始字符串;

    2)分割字段,提取需要参与统计的列(如一级分类 prod_cat、均价 avg_price);

    3)输出中间键值对:

key=分组字段,value=待聚合数值。

  1. 实训示例:

一行数据 1,蔬菜,白菜,1.0→ map 提取输出 key:蔬菜,value:1.0。

  1. 特点:多台机器上的 Map 任务完全并行、互不干扰,只处理自己分到的数据分片。

阶段 3:Shuffle 洗牌阶段(MapReduce 自动执行,不用手写代码)

这是连接 Map 和 Reduce 的核心桥梁,也是分布式分组的关键,分为 5 个小步骤:分区、排序、合并、分组、远程拉取。

  1. 分区 Partition:根据 Map 输出的 key 划分数据,相同 key 的数据会被分配到同一个 Reduce 任务;
  2. 本地排序 Sort:每台机器 Map 输出的中间数据自动按 key 升序排序;
  3. 合并 Combine(可选):本地提前聚合同 key 数据,减少网络传输;
  4. 分组 Group:把所有相同 key 的 value 归为一个集合;
  5. 远程复制 Fetch:通过网络,把不同机器上相同 key 的数据全部拉取到同一台 Reduce 服务器。
  6. 实训效果:所有蔬菜的价格、所有水果的价格自动汇总到一起,等待 Reduce 计算平均值。

阶段 4:Reduce 聚合阶段(reduce () 方法)

  1. 核心功能:接收 Shuffle 分组后的一组同 key 数据,执行聚合计算(求和、平均、计数、最大值、最小值等)。

  2. 具体完成工作:

    1)接收一组数据:key=蔬菜,values=[1.0,2.1];

    2)遍历该 key 下全部 value,完成聚合逻辑(求和、求平均、去重、统计条数);

    3)输出最终业务结果键值对

  3. 实训示例:

    key = 蔬菜,价格集合 [1.0,2.1] → reduce 计算平均 1.55,输出 蔬菜:1.55。

阶段 5:Output 结果输出阶段

  1. 完成工作
  2. 将 Reduce 计算完成的最终结果写入 HDFS;
  3. 每个 Reduce 任务单独生成一个输出文件 part-r-0000x;
  4. 数据持久化保存,可供 Hive 查询、Java 程序同步至 MySQL、可视化大屏读取。

1782375043917

其中mapreduce慢就慢在shuffle了,map得到的kv结果会存放到100M内存缓冲区,超过80%会写到磁盘,同时会对文件分组(reduce个数取模的个数相同组数)内部排序,生成多个有序的小文件,小文件不便于传输,于是小文件进行合并为大分组有序的文件。

接下来reduce从各个map中拉取所需要的数据文件过来,拉取到缓存,也会写磁盘,生成小文件,然后对小文件进行合并大文件,大文件也要按同k分组按k排序,reduce对这文件进行value值处理。

整个过程大量落盘,map阶段内存缓存区达到80M落盘,合并成大文件又落盘了。reduce阶段拉取文件进入内存缓存区,达到阈值又小文件落盘,小文件合并大文件又合并,然后交给reduce处理。

mapreduce 不消耗内存,大量磁盘交互,在大文件在移动到reduce过程走网络,效率也低

三、结合实训案例对应关系(农产品价格统计)

需求:统计每种一级分类的平均价格,等价 HQL:

SELECT prod_cat, ROUND(AVG(avg_price),2) FROM ods_xinfadi_price GROUP BY prod_cat;
  1. Input:读取 HDFS 中所有新发地 CSV 原始文件,切分数据分片;
  2. Map:每行拆分字段,提取分类 + 均价,输出 <分类,价格>;
  3. Shuffle:自动把所有 “蔬菜” 价格、“水果” 价格分别归集;
  4. Reduce:同一品类所有价格求和 / 计数,算出平均价;
  5. Output:把品类 + 均价结果写入 HDFS,存入 ADS 指标层供大屏使用。

四、Map、Reduce 各自独立分工总结

Map 职责(处理单行明细)

  • 数据读取、字段提取;
  • 脏数据过滤、格式转换;
  • 产出分组用的中间键值对;
  • 并行处理,适合海量明细遍历。

Reduce 职责(分组聚合)

  • 接收同一分组的全部数据;
  • 聚合运算:AVG、SUM、COUNT、MAX、MIN;
  • 生成最终统计指标;
  • 一组 key 只执行一次 reduce 逻辑。

五、MapReduce 整体作用总结(可直接写入实训报告)

  1. 实现海量离线数据分布式并行计算,突破单机处理数据量上限;
  2. 分离数据映射(Map)与聚合统计(Reduce),降低大数据开发难度;
  3. 内置 Shuffle 自动完成数据分组、排序、分发,开发者无需手动处理多服务器数据传输;
  4. 是 Hive 数仓底层执行引擎,所有 HQL 分组、聚合、筛选语句都会翻译为 MapReduce 任务运行;
  5. 计算结果持久化到 HDFS,支撑后续指标同步、数据可视化上层业务。

上一篇
2_hadoop →