MapReduce
约 1866 字大约 6 分钟
2026-09-12
MapReduce虽然功能上和GFS完全不一样,但是设计思想有不少相同之处,也是分布式的经典。
模型
MapReduce的输入是一个kv键值对集合。用户负责实现Map函数和Reduce函数。Map函数读取输入,并输出中间kv集合。Reduce函数负责将一个key和其对应的一些value读入,然后将这些value合并为一个可能更小的value集合。中间数据通过迭代器传入,这样我们就不需要超大内存来处理了。
范例:word count,分布式grep,倒排索引等
环境:2-4 GB内存 x86 Linux机器集群,存储方式为GFS+低成本硬盘。
实现
执行过程
将输入文件切分为M片,平均每片16到64 MB。
master将任务分配给worker,可能是map任务也可能是reduce任务。
执行map任务的worker读取分片的内容,解析出kv,并传递给Map函数。Map过程中产生的kv会被放在内存中。
Map产生的kv会周期性地被写入磁盘,被切分为R个区域。这些kv在磁盘的地址会被传回master,以将他们转发给reduce worker。
reduce用RPC读取这些缓冲的数据,并根据key来排序。由于key很多,所以不排序内存会不够用(我们需要让同key数据同批次执行),并且我们需要使用外排序。
reduce遍历有序数据,对于每个遇到的唯一中间键,系统会将该key及对应的value集合传递给用户定义的Reduce函数。
通常来说,用户不需要将最后的输出结果合并,而是将它们转发到下一个MapReduce调用。
Master
如前文所述,master会存储缓冲数据的位置和大小,并且将这些信息的更新视作map任务的完成。除此之外,它还存储worker的状态(idle,in-progress,completed)和身份信息。
容错
由于MapReduce处理的是非常大量的数据,所以和GFS一样,它也必须容忍机器的故障。
master会周期性地ping每个worker,如果一段时间没有回复,就会将该worker标记为失败。任何一个map任务在被完成之后都会被重置为idle状态,失败的也一样,这样它们才能够被调度到别的worker上。已完成的map任务在发生故障时会重新执行,因为其输出结果存储在故障机器的本地磁盘上,因而无法访问;而已完成的reduce任务则无需重新执行,因为其输出结果存储在全局文件系统中。如果一个map任务由于故障被换到了新worker上执行并通知其它所有执行reduce任务的worker,这样它们就会从新worker上读取map任务的数据。所以,MapReduce具备大规模worker故障的恢复能力。
将master的数据备份十分容易,但是由于master的故障概率十分小,所以我们一般在故障时让客户端直接重新开始MapReduce任务。
由于这是分布式应用,所以我们要花点心思来保证它的执行结果与顺序执行是完全一致的。所有的任务都会将其结果写入一个私有临时文件。当map任务完成时,worker会向master发送其生成的R个临时文件的名字。如果master收到重复的完成信息,就会忽略,否则它会记录这R(reduce的worker数)个文件的名字。当一个reduce任务完成后,worker会原子地讲它的临时输出文件重命名为最终输出文件的名字。我们依赖底层文件系统提供的原子重命名操作,以确保最终的文件系统状态仅包含由单次reduce任务执行所产生的数据。
局部性
网络带宽资源十分紧张,所以我们用GFS来节省。GFS对同一份文件有三个或以上数量的副本,MapReduce的worker就会尽量使用离他更近的那个副本。对于GFS不在赘述,可见之前的文章。
任务粒度
我们将map阶段分为M片,reduce阶段分为R片。理想情况下M和R都要远远大于worker的数量,这样有助于负载均衡。通常情况下时间复杂度为 O(M+R),空间复杂度为 O(M∗R)。我们通常选择M使得每个独立任务的输入数据量大致在16 MB至64 MB之间,并将R设为预期使用的worker数量的一个小整数倍。我们经常采用M = 200000且R = 5000的参数进行MapReduce计算,此时使用2000台worker。
备份任务
对MapReduce的完成时间影响最大的就是花费时间特别长的一小部份任务,所以我们需要一个机制来缓解这些“滞留者”。当一个MapReduce任务快要完成时,master会调度对剩余in-progress任务进行备份。无论原来的或备份的任务被完成,这个任务都会被标记为完成。我们对这一机制进行了优化,使其通常将操作所使用的计算资源增加幅度控制在不超过百分之几范围内,事实证明十分有效(事实见原论文)。
改进措施
切分函数
我们通常用哈希函数来切分以保证负载均衡,但有些情况下我们可能需要单独设计哈希函数。比如key是URL的情况,我们可能会希望同一主机的URL分配到一个worker去运算,此时哈希函数为 hash(Hostname(urlkey))。
顺序保证
保证在每一个分区中的中间kv都是有序的,这样更方便我们有序的输出,对某些业务十分友好。
组合器函数
以word count任务为例,通常来说词频符合Zipf distribution,所以每个map任务都会产生大量的重复的类似 <the, 1> 的kv,我们允许用户自定义组合器函数。组合器函数在执行map任务的机器上被执行,逻辑通常和Reduce函数是一样的,只是组合器函数的输出结果会被写入一个中间文件,该文件随后将被发送至reduce任务。
跳过异常的记录
有时候代码中的bug难以修复,而某些情况下忽略少量记录又是可以接受的(这在搜广推业务很常见),所以MapReduce有一个选项允许检测某些异常记录并跳过它们。在Map和Reduce被唤醒前,MapReduce会保存参数的序列号。如果用户异常,worker会将序列号发送给master。如果该序列号出现多个异常,它就会被跳过直到被重新执行。
展望
和GFS一样,MapReduce有太多缺点了:不支持流式计算、静态集群、单一master等等。但它已然成为分布式计算的“汇编语言”。后来的Spark、Flink等应用都是对MapReduce在不同维度的抽象进行发展和重构。
