Apache Hadoop — 在一个 Pod 里搭起并运维 HDFS 与 YARN
MapReduce 的开销集中在 Map 与 Reduce 之间的 Shuffle
一句话总结
MapReduce 是 Map → (分区、排序、shuffle)→ Reduce 的框架,用户只需要编写前后两个函数。中间的排序和 shuffle 由框架完成,作业的大部分开销都出在那里。Combiner 是缩减中间部分的装置,计数器是窥视中间部分的窗口。
为什么需要这样的形态
假设要在分散于几百台机器上的 TB 级数据中统计单词出现的次数。把数据汇集到一台服务器上,网络会最先垮掉。这就是 HDFS 设计文档提出的原则“把计算移到数据旁边,比移动数据更便宜”的由来。这样一来,就得在块所在的每台服务器上分别运行计算,而这次又必须把分别运行的结果按相同的单词重新汇集起来。服务器挂了,就得把它那一份重新运行,慢的服务器则要等待。
如果每个作业都重新解决这个问题,代码的大部分就都成了分布式处理的善后工作。MapReduce 把这些善后工作交给了框架。MapReduce 教程总结道:框架负责任务的调度和监控,以及失败任务的重新执行,应用只需提供输入输出位置和 map、reduce 函数。代价是所有计算都必须表示成键值对。有了这个限制,“相同的键会去往同一个 Reducer”这一个约定,就解决了汇集的问题。
工作原理
把流程写成键值对的形式是这样的,和教程中的完全一致。
(입력) <k1, v1> -> map -> <k2, v2> -> combine -> <k2, v2> -> reduce -> <k3, v3> (출력)
Map 端。根据教程,Map 的数量通常取决于输入的总大小,也就是输入文件的块数。一个块通常成为一个输入分片(split),每个分片会启动一个 Map 任务。Map 产生的键值对不会立刻写入磁盘,而是先堆积在内存缓冲区里。在 mapred-default.xml 中,这个缓冲区大小 mapreduce.task.io.sort.mb 的默认值是 100MB(任务堆只有 180MB 的实验镜像把它缩小到了 32MB),达到 mapreduce.map.sort.spill.percent 的默认值 0.80 时,后台线程会把内容排序后溢写(spill)到磁盘。Map 结束后,会把溢写出的各个文件合并成一个。这个文件按 Reducer 的数量被划分成若干分区,每个分区内部按键排好序。
哪个键去往哪个分区,由分区器(Partitioner)决定。默认的是 HashPartitioner,即键的哈希值对 Reducer 数量取余。相同的键,无论来自哪个 Map,都会去往编号相同的 Reducer。汇集之所以成立,依据就是这一行。
Reduce 端。教程说明,Reducer 要经过 shuffle、排序、reduce 三个阶段。在 shuffle 中,它只通过 HTTP 取回所有 Map 输出中属于自己的分区,取回的同时做归并排序,把相同键的值排成一行。然后对每个键调用一次 reduce(키, 값 목록)(占位符依次为键与值列表)。如果把 Reducer 的数量设为 0,Map 的输出就直接成为结果,也不做排序。
Reducer 数量 mapreduce.job.reduces 的默认值是 1。不设置的话,全部 Map 输出都会涌向同一个 Reducer,结果文件也只有一个 part-r-00000。把 Reducer 增加到 3 个,结果文件就变成 3 个。Reducer 是按键的顺序被调用的,所以像单词计数这种直接使用收到的键的作业,每个文件内部是按键排序的。但正如教程所强调的,框架不会再对 Reducer 的输出重新排序,把三个文件首尾相接,整体的顺序是乱的。因为它是按哈希划分的。
要想整体排成一行,就得更换分区器。TeraSort 示例展示了这个方法。从输入中抽取键的样本,确定 N-1 个边界,每个边界区间分配一个 Reducer。这样,第 i 号 Reducer 的输出全都小于第 i+1 号,把文件按编号顺序首尾相接,整体就是有序的。TeraValidate 就是用来检查这一点的作业。
Combiner 与计数器
单词计数的 Map 会输出几万次 (the, 1)。如果原封不动地 shuffle,网络上就要传几万个 1。Combiner 是在 Map 端事先缩减成 (the, 38211) 再发送的本地聚合。教程指出,Combiner 会在溢写时运行,在 Reduce 端合并期间也可能运行,而大于缓冲区的记录是否经过 Combiner,并没有规定。总结起来就是:Combiner 会运行几次是没有保证的。因此,只有运行一次或运行三次结果都相同的运算,才能用作 Combiner。求和与最大值可以,平均值不行——平均值的平均值不是平均值。
Combiner 是否发挥了作用,要看计数器。作业结束后,会打印出框架统计的值,例如 Map 输入记录和 Map 输出记录、Combiner 输入和输出记录、Reduce 输入分组数、溢写记录数。如果 Combiner 的输出比输入小得多,就说明 shuffle 也相应减少了。如果是单词计数,Reduce 输入分组数就是结果中不同单词的个数。计数器是不必翻日志,就能用数字证明作业做了什么的最便宜的办法。
运行在 YARN 之上意味着什么
自 Hadoop 2 之后,MapReduce 不再自己管理资源。正如 YARN 架构文档所述,职责分成了协调整个集群的 ResourceManager,和每个应用各启动一个的 ApplicationMaster,MapReduce 作业的 ApplicationMaster 就是 MRAppMaster。提交一个作业,至少会启动一个 AM 容器,再加上与 Map 任务数量相同的容器,以及与 Reducer 数量相同的容器。
作业结束后,AM 会留下历史记录。默认的暂存目录是 /tmp/hadoop-yarn/staging,中间完成目录是其下的 history/done_intermediate。即使不启动作业历史服务器,这里也会按用户累积 .jhist(事件历史)、.summary 和作业配置 XML。历史服务器只负责转移并展示这些文件,记录本身是由 AM 写的。
在现场相遇的样子
第一,作业不是在集群上,而是在本地运行。mapreduce.framework.name 的默认值是 local。如果客户端没有改成 yarn,从它提交的作业不会出现在 ResourceManager 中,而是在提交作业的那台机器的一个 JVM 里悄悄运行。对“为什么 ResourceManager 页面上没有作业”这个问题,这是最常见的答案。实验镜像已经把这个值改成了 yarn。
第二,没有 AM 可以放的位置。yarn.app.mapreduce.am.resource.mb 的默认值是 1536MB。在 NodeManager 分给容器的份额比它小的集群里,没有哪个节点能放下 AM 容器,作业连启动都做不到。在这个实验环境这样份额只有 1024MB 的地方,必须另外调低 AM 的大小,作业才能运行。实验镜像把 AM 设为 384MB,Map 和 Reduce 设为 256MB,让一个 AM 和两个任务可以同时放得下。
第三,有一个 Reducer 一直结束不了。要么是一直保持默认值 1,要么是某一个键上堆积了大量的值。前者靠配置来修正,后者靠键的设计来修正。
第四,因为 Combiner 导致数字不对。如果把平均值或中位数的 Reducer 代码原样挂成 Combiner,结果会随 Combiner 的运行次数而变化。这种时对时错的 bug 很难找。
实际工作中真正重要的事
- 相同的键会去往同一个 Reducer。汇集的依据就是分区器的那一行。
- shuffle 就是开销。减少 Map 和 Reduce 之间传输的字节,是调优的大部分内容。
- Combiner 会运行几次是不确定的。只使用应用多次结果也相同的运算。
- Reducer 的默认值是 1 个。结果文件数和并行度都取决于这个值。
- framework.name 的默认值是 local。如果作业不在 ResourceManager 中,先看这个。
- 用计数器来证明。输入、输出、Combiner 和分组数,用数字展示了作业所做的事。
下一项实验要做什么
先启动 YARN,确认 NodeManager 已注册,然后把五本书上传到 HDFS,运行示例 jar 中的单词计数作业,并记下作业 ID。在没有历史服务器的情况下,找到 HDFS 的 done_intermediate 中留下的作业历史文件并打开,实验镜像已让这个文件以人可读的 JSON 行来写入。从作业结束时的全部计数器中取出 Map 输入与输出、Combiner 输入与输出、Reduce 输入分组数和输出记录,计算 Combiner 把 shuffle 缩减了多少。把 Reducer 增加到三个,看结果文件是如何划分的,最后用 teragen、terasort、teravalidate 运行全局排序,并通过验证。