目录

3.1 剖析MapReduce作业运行机制

(1)提交 MapReduce 作业

(2)初始化 MapReduce 作业

(3)分配任务

(4)执行任务

(5)更新进度和状态

3.2 作业失败与容错

(1)任务容错

(2)application master容错

(3)NodeManager容错

(4)ResourceManager容错

3.3 shuffle过程详解

(1)Map 端

(2)Reduce 端

3.1 剖析MapReduce作业运行机制

通过方法调用运行MapReduce作业有两种方式,一种是job对象的waitForCompletion()方法,另一种是Job对象的sumit()方法。

waitForCompletion()方法底层也调用了sumit()方法,它用于提交以前没有提交过的作业,并等待作业运行完成。sumit()方法调用封装了大量的处理细节,它仅仅提交作业即可,无需等待作业的完成。

MapReduce作业的运行6原理如图所示:

按照大模块来划分,图 中一共包含5个独立的实体。

  • 客户端:负责提交MapReduce作业。

  • YARN资源管理器(ResourceManager) :负责协调集群上计算机资源的分配。

  • YARN节点管理器(NodeManager) :负责启动和监视集群中机器上的计算容器(Container)

  • MapReduce的Application Master:负责协调运行MapReduce作业的任务。它与MapReduce任务都运行在Container中,这些容器由ResourceManager分配并由NodeManager管理。

  • 分布式文件系统:一般为HDFS,用来与其他实体共享作业文件。

MapReduce作业完整的运行过程如下。

(1)提交 MapReduce 作业

通过调用 Job 对象的 waitForCompletion()方法来运行MapReduce 作业, waitForCompletion()的底层调用submit0方法(步骤1) 。在submit()方法的内部创建了一个JobSummiter实例,并且调用其submitJobInternal()方法。提交MapReduce作业后, waitForCompletion()每秒会轮询作业的进度,如果发现自上次报告后有改变,则会将进度报告打印到控制台。MapReduce作业完成后,如果运行成功,就显示作业计数器;如果运行失败,则将作业失败的错误日志打印到控制台。

JobSummiter所实现的作业提交过程如下。

1)向ResourceManager请求一个新应用ID,该ID设置为MapReduce作业的ID(步骤2) 。

2)检查MapReduce作业的输出路径。例如,如果没有指定输出目录或输出目录已经存在,它会拒绝提交作业,同时将错误信息反馈给MapReduce程序。

3)计算MapReduce作业的输入分片,如果分片无法计算,例如,输入路径不存在,就会都绝提交作业,同时将错误信息反馈给MapReduce程序。

4)将作业运行所需要的资源(包括作业jar文件、配置文件和计算得到的输入分片)复制到共享文件系统中,存储在以作业ID命名的目录下(步骤3)。

5)通过ResourceManager来提交作业(步骤4) 。

(2)初始化 MapReduce 作业

ResourceManager收到调用它的submitApplication()消息后,便将请求传递给YARN调度器(Scheduler) 。Scheduler分配一个容器,然后ResourceManager 在 NodeManager 的管理下在容器中启动Application Master 进程(步骤5a和5b)。

MapReduce作业的Application Master是一个Java应用程序,它的主类是MRAppMaster。Application Master会对作业进行初始化(步骤6),通过创建多个簿记对象来跟踪作业的执行进度和完成报告。接着Application Master接收来自共享文件系统的、在客户端计算的输入分片(步骤 7)。然后对每一个分片创建一个 Map 任务对象以及通过作业设置的多个 Reduce 任务对象任务 ID 会在此时分配。

(3)分配任务

Application Master 为作业中的Map 任务和 Reduce任务向 ResourceManager 请求容器(步骤8)。它首先为Map任务发出请求,该请求优先级要高于Reduce任务的请求,这是因为所有的Map任务必须在Reduce的排序阶段能够启动前完成,直到有5%的Map任务已经完成时,为Roduce任务的请求才会发出。Reduce任务可以在集群任意位置运行,而Map任务需要遵循数据本地化原则,选择就近的节点运行。

Application Master发送的请求也为任务指定了内存大小和CPU数量,默认情况下,每个Map任务和Reduce任务都分配到1GB的内存和一个虚拟的内核,当然这些值都可以通过参数进行配置。

(4)执行任务

一旦ResourceManager的调度器为任务分配了一个特定节点上的容器,Application Master就通过与 NodeManager 通信来启动容器(步骤 9a 和 9b)。该任务由主类为 YamChild (在指定的JVM中运行)的一个Java应用程序执行。在它运行任务之前,首先将任务需要的资源本地化,包括作业的配置、jar文件和所有来自分布式缓存的文件(步骤10) ,然后运行Map任务和Reduce 任务(步骤 11)。

(5)更新进度和状态

MapReduce作业是长时间运行的批量作业,运行时间从数秒到数小时,所以对于用户来说,能够知晓关于作业运行进度的一些反馈信息非常重要。一个作业和它的任务都有一个状态(Status) ,包括作业或任务的状态(如运行中、成功、失败)、Map和Reduce的进度、作业计任数器的值等。

在Map任务和Reduce任务运行时,子进程和自己的父Application Master通过接口进行通信,默认每隔3s,任务通过这个接口向自己的Application Master报告进度和状态(包括计数器), Application Master会形成一个作业的汇聚视图。

ResourceManager的 Web 界面显示了所有运行中的应用程序,并且分别有链接指向这些应用各自的Application Master界面,这些界面展示了MapReduce作业的更多细节,包含其进度。

在作业运行期间,客户端每秒轮询一次 Application Master,从而获取最新的状态。客户端也可以使用Job的getStatus()方法得到一个JobStatus的实例,该实例包含作业的所有状态信息。

状态更新在 MapReduce 系统中的传递流程如图所示:

(6)完成作业

当Application Master收到作业最后一个任务已完成的通知后,就会将作业的状态设置为“成功”。在Job轮询状态时,就会知道任务已经成功完成,于是就会打印一条信息告知用户,然后从 waitForCompletion()方法返回。Job 的统计信息和计数值也在这个时候输出到控制台。

MapReduce 作业完成时,Application Master 和任务容器会清理其工作状态。Hadoop 集群会存储作业信息,方便用户随时查询。

3.2 作业失败与容错

在MapReduce作业运行过程中,可能会出现各种问题,例如,用户代码错误、作业进程崩溃、机器故障等。使用Hadoop最主要的好处之一就是它能自动处理这些故障,让用户能够成功运行作业。在作业运行过程中,需要考虑的容错实体包括任务、Application Master、 NodeManager和ResourceManager。

(1)任务容错

当application master被告知一个任务尝试失败后,它将重新调度该任务的执行。application master会试图避免在之前失败过的NodeManager上重新调度该任务。此外,如果一个任务失败数超过4次,该任务将不会再尝试执行。

(2)application master容错

application master向ResourceManager发送周期性的心跳,当application master失败时,ResourceManager将检测到该失败,并在一个新的容器中重新启动一个application master实例。对于新的application master来说,它将使用作业历史记录来恢复失败的应用程序所运行任务的状态,所以这些任务不需要重新运行。

(3)NodeManager容错

如果一个NodeManager节点因中断或运行缓慢而失败,那么它就会停止向ResourceManager发送心跳信息(或者发送频率很低)。默认情况下,如果ResourceManager在10分钟内没有收到一个心跳信息,它将会通知停止发送心跳信息的NodeManager,并且将其从自己的节点池中移除。

在出现故障的NodeManager节点上运行的任何任务或application master,将会按前面描述的机制进行恢复。另外,对于出现故障的NodeManager节点,那么曾经在其上运行且成功完成的map任务,如果属于未完成的作业,那么application master会安排它们重新运行。这是因为它们的中间输出结果是存放在故障NodeManager节点所在的本地文件系统中,reduce任务可能无法访问。

(4)ResourceManager容错

ResourceManager出现故障是比较严重的,因为没有ResourceManager,作业和任务容器将无法启动。在默认的配置中,ResourceManager是一个单点故障,因为在机器出现故障时,所有的作业都会失败并且不能被恢复。

为了实现高可用(HA),有必要以一种active-standby配置模式运行一对ResourceManager。如果active ResourceManager出现故障,则standby ResourceManager可以很快的接管,并且对客户端来说没有明显的中断现象。

3.3 shuffle过程详解

MapReduce确保每个Reducer的输入都是按key排序的,系统执行排序、将Map输出作为输入传给Reducer的过程称为Shuffle。接下来将重点介绍Shuffle 是如何工作的,这样有助于理解 MapReduce 工作机制。MapReduce 的 Shuffle 过程如图所示。

(1)Map 端

Map任务开始输出中间结果时,并不是直接写入磁盘,而是利用缓冲的方式写入内存,并出于效率的考虑对输出结果进行预排序。

每个Map任务都有一个环形内存缓冲区,用于存储任务输出结果。默认情况下,缓冲区的大小为100MB,这个值可以通过mapreduce.task.io.sort.mb属性来设置。

一旦缓冲区中的数据达到阈值(默认为缓冲区大小的80%) ,后台线程就开始将数据刷写到磁盘。在数据刷写到磁盘的过程中,Map任务的输出将继续写到缓冲区,但是如果在此期间缓冲区被写满了,那么Map会被阻塞,直到写磁盘过程完成为止。

在缓冲区数据刷写到磁盘之前,后台线程首先会根据数据被发送到的Reducer个数,将数据划分成不同的分区(Partition)。在每个分区中,后台线程按照key在内存中排序,如果此时有一个 combiner 函数,它会在排序后的输出上运行。运行 combiner 函数可以减少写到磁盘和传递到 Reducer的数据量。

每次内存缓冲区达到溢出阈值,就会刷写一个溢出文件,当 Map 任务输出最后一条记录之后会有多个溢出文件。在Map 任务完成之前,溢出文件被合并成一个已分区且已排序的输出文件。默认如果至少存在3个溢出文件,那么输出文件写到磁盘之前会再次运行Combiner函数。如果少于 3 个溢出文件,则不会运行 Combiner 函数,因为 Map 输出规模太小不值得调用Combiner 函数(带来的开销较大)。在 Map 输出写到磁盘的过程中,还可以对输出数据进行压缩,加快磁盘写入速度,节约磁盘空间,同时也减少了发送给Reducer的数据量。

(2)Reduce 端

Map输出文件位于运行Map任务的NodeManager的本地磁盘,现在NodeManager需要为分区文件运行Reduce任务,而且 Reduce任务需要集群上若干个Map任务的Map输出作为其特殊的分区文件。每个Map任务的完成时间可能不同,因此在每个任务完成时,Reduce任务就开始复制其输出。这就是 Reduce 任务的复制阶段。默认情况下,Reduce 任务有 5 个复制线程,因此可以并行获取 Map 输出。

如果 Map 输出结果比较小,数据会被复制到 Reduce 任务的 JVM 内存中;否则,Map 输出会被复制到磁盘中。一旦内存缓冲区达到阈值大小,数据合并后会刷写到磁盘。如果指定了Combiner函数,在合并期间可以运行Combiner函数,从而减少写入磁盘的数据量。随着磁盘上溢出文件的增多,后台线程会将它们合并为更大的、已排序的文件,这样可以为后续的合并节省时间。

复制完所有 Map 输出后,Reduce任务进入排序阶段,这个阶段将 Map 输出进行合并,保持其顺序排序。这个过程是循环进行的,例如,如果有50个Map输出,默认合并因子为10,那么需要进行5次合并,每次将10个文件合并为一个大文件,因此最后有5个中间文件。

在最后的 Reduce 阶段,直接把数据输入 reduce 函数,从而节省了一次磁盘往返过程。因为最后一次合并并没有将这5个中间文件合并成一个已排序的大文件,而是直接合并到Reduce作为数据输入。在Reduce阶段,对已排序数据中的每个key调用reduce函数进行处理,其输出结果直接写出到文件系统,这里的文件系统一般为 HDFS。

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐