1.1 什么是 MapReduce


文档摘要

1.1 什么是 MapReduce 本节摘要:MapReduce 是 Google 于 2004 年公开发表、随后由 Hadoop 开源落地的分布式计算模型。它把"分而治之"抽象成 map 与 reduce 两个用户函数,把并行、调度、容错、数据分发全部收编为框架职责。本节从定义出发,还原它诞生时的工程命题,拆开它的两函数契约,最后用一个可以在本机复算的迷你例子把"模型"落到可执行代码上。 一句话定义与它的三层含义 MapReduce 是一种用于大规模数据集并行计算的编程模型:用户只需编写 map 和 reduce 两个函数,框架负责把作业切分、分发到成百上千台机器上执行,并处理失败重试与机器间的数据传输。这句话可以再剥成三层: 第一层,它是模型而非产品。

1.1 什么是 MapReduce

本节摘要:MapReduce 是 Google 于 2004 年公开发表、随后由 Hadoop 开源落地的分布式计算模型。它把"分而治之"抽象成 map 与 reduce 两个用户函数,把并行、调度、容错、数据分发全部收编为框架职责。本节从定义出发,还原它诞生时的工程命题,拆开它的两函数契约,最后用一个可以在本机复算的迷你例子把"模型"落到可执行代码上。

一句话定义与它的三层含义

MapReduce 是一种用于大规模数据集并行计算的编程模型:用户只需编写 map 和 reduce 两个函数,框架负责把作业切分、分发到成百上千台机器上执行,并处理失败重试与机器间的数据传输。这句话可以再剥成三层:

第一层,它是模型而非产品。map 与 reduce 借自 Lisp 等函数式语言的列表操作:map 对每条记录独立变换,reduce 把同键的值归并。模型不绑定任何实现——Google 的内部 C++ 实现、Hadoop 的 Java 实现、乃至今天在教材里用单机模拟的教学实现,跑的都是同一套语义。

第二层,它是面向普通工程师的分布式系统封装。2004 年之前,写一个并行的日志统计程序,工程师要自己处理数据划分、任务分发、机器故障、慢机器拖尾、结果汇总,代码里业务逻辑常常只占三成。MapReduce 的革命性在于:这七成"分布式脏活"被收进框架,业务代码收敛到两个函数签名里。

第三层,它的运行单位是作业。一个 MapReduce 作业把输入数据切分成若干分片,每个分片由一个 Map 任务处理;Map 的中间输出按键分区后,被汇聚到若干 Reduce 任务做最终归并。作业、任务、分片、分区这四个词,构成了后面所有章节的词汇表。

层次 关注点 谁负责
模型 map 与 reduce 的输入输出契约 用户
框架 切分、调度、Shuffle、容错、计数 实现
作业 一次计算从提交到产出的生命周期 两者协作

诞生的工程命题:索引构建倒逼出的范式

理解一个技术"为什么出现在那个时间点",比背它的定义更有用。2003 到 2004 年,Google 面对的现实是:网页规模冲到数十亿页,为了给搜索引擎建倒排索引,要把全网页面解析、提词、统计词频、再按词聚合一轮。这个计算量单机无解,只能横跨上千台廉价机器。

廉价机器意味着每月必有若干台损坏,于是任何并行程序都必须回答三个问题:某台机器中途挂了,它已完成的部分怎么办?某台机器特别慢,整个作业要不要等它?数据分布在各台机器的磁盘上,怎么把它们搬到一起做聚合?在 MapReduce 之前,每个并行程序各自 answering,答案还常常是错的。

Google 的观察很朴素:内部大量批处理计算——倒排索引、排序、词频统计、链接图分析——都能套进同一个壳子:"对每条记录独立映射,再按键归并"。既然模式重复出现,就值得把这个壳子做成基础设施。2004 年的论文报告了彼时的成绩:单次索引构建从数千行手写并行代码,缩到几百行 MapReduce 代码;上千节点的作业自动挺过机器故障。次年 Hadoop 把这套思想开源成 Java 实现,MapReduce 从一家公司的内部武器变成行业公共底座。

论文前:并行逻辑 + 容错 + 数据搬运 + 业务 约 3000 行 论文后:map 函数 + reduce 函数 + 作业配置 约 300 行 省下的 90% 代码,全部长进了框架

两函数契约:模型的最小接口

把模型压缩到最小,就是两个函数签名(以 Hadoop Java API 为例):

// map:输入一对,输出零到多对中间键值 protected void map(KEYIN key, VALUEIN value, Context context) { // 对 value 做切分、过滤、变换 context.write(中间键, 中间值); } // reduce:输入一个键和它的全部值,输出零到多对最终结果 protected void reduce(KEYIN key, Iterable<VALUEIN> values, Context context) { // 遍历 values 做求和、最值、拼接 context.write(key, 最终值); }

契约的关键约束只有两条:map 对每条输入必须互相独立(不依赖处理顺序、不依赖其他记录),reduce 只通过键拿到它该看的数据(不同键之间不通信)。只要业务能接受这两条,框架就能把 map 随便复制到任何节点、把 reduce 的输入提前排好序送上门。反过来,一旦你在 map 里偷偷记住了"上一条记录",或是在 reduce 里想访问另一个键的结果,范式就破了——这也是第 7 章讨论其局限的伏笔。

在两函数之间,框架自动补齐了一段用户不写却真实存在的逻辑:map 的输出按中间键分区,同键数据跨节点聚到同一个 reduce,且送达时按键有序。这段"数据迁徙"就是 Shuffle,第三幕的主角。

图 1.1-1 模型全貌:两函数与框架分工

图中实线是数据流向,虚线标注职责归属:用户只拥有两个函数,其余整条流水线都姓"框架"。

一个可以亲手复算的迷你例子

光说不练假把式。下面用原生 Java MapReduce 写一个"按部门汇总工资"——它是 WordCount 的近亲,但带了一点业务味。输入是若干行"姓名,部门,工资":

赵敏,研发,12000 钱磊,研发,15000 孙蕙,市场,9000 李昂,市场,11000 周婷,研发,13000

map 阶段把每行拆开,以部门为键、工资为值原样发出;reduce 阶段把同部门的工资累加:

public class DeptSalarySum { public static class DeptMapper extends Mapper<LongWritable, Text, Text, LongWritable> { private Text dept = new Text(); private LongWritable salary = new LongWritable(); @Override protected void map(LongWritable offset, Text line, Context context) { String[] cols = line.toString().split(","); if (cols.length != 3) return; // 脏行直接跳过 dept.set(cols[1]); salary.set(Long.parseLong(cols[2])); context.write(dept, salary); } } public static class SumReducer extends Reducer<Text, LongWritable, Text, LongWritable> { private LongWritable result = new LongWritable(); @Override protected void reduce(Text dept, Iterable<LongWritable> salaries, Context context) throws IOException, InterruptedException { long sum = 0; for (LongWritable s : salaries) sum += s.get(); result.set(sum); context.write(dept, result); } } public static void main(String[] args) throws Exception { Job job = Job.getInstance(); job.setJarByClass(DeptSalarySum.class); job.setMapperClass(DeptMapper.class); job.setReducerClass(SumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(LongWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }

提交运行:

hadoop fs -put emp.csv /input/dept/ hadoop jar dept-sum.jar DeptSalarySum /input/dept /output/dept hadoop fs -cat /output/dept/part-r-00000

输出(五条记录、两个部门,可手工验算):

市场 20000 研发 40000

值得在脑内推演一遍框架在这五行数据上做了什么:假设输入被切成两个分片,前两行与后三行各归一个 Map 任务。第一个 Map 发出(研发,12000)(研发,15000),第二个发出(市场,9000)(市场,11000)(研发,13000);框架按"部门"分区,市场组进一个 Reduce,研发组进另一个 Reduce,且各自按键排好序送达。于是 reduce 函数拿到的 values 是有序列表,求和逻辑水到渠成。数据怎么切、怎么搬、怎么排序,用户代码里一个字没提——这就是模型的意义。

与相邻概念的边界

初学者常把 MapReduce 与三个近邻混为一谈,这里划清边界:

概念 关系 一句话区分
HDFS 共生不隶属 MapReduce 是计算,HDFS 是存储,前者读后者的块
Hadoop 实现与其上的生态 MapReduce 是 Hadoop 最早的计算引擎之一
Spark 后来者 同一范式思想(宽依赖处发生 Shuffle),引擎换成内存优先

还有一对易混词:数据并行任务并行。MapReduce 是数据并行——所有 Map 任务跑的是同一段代码,只是各管一份数据;它不是把一个算法的不同步骤分给不同机器。认清这一点,就能解释为什么矩阵乘法这种需要全局通信的算法在 MapReduce 里写起来格外别扭,而日志统计类任务写起来行云流水。

本节自查

  1. 用自己的话说出 map 与 reduce 各自的输入输出,中间缺的那段逻辑由谁补;
  2. 复述 2004 年 Google 面对的三个分布式难题,说明框架分别用什么机制接住;
  3. 把上面迷你例子的输入再添两行,手工推演两个 Map 任务各自发出哪些中间对;
  4. 判断"计算每个用户的相邻两笔交易时间差"是否适合该范式,说出依据(提示:相邻性破坏了 map 的独立性)。

下一节把镜头拉远:这个两函数模型作为"范式",究竟强在哪、贵在哪。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U