MapReduce 论文笔记
起源
几位作者在工作中针对巨量的数据,写了上百种数据处理程序,它们有一些特征:
- 业务逻辑简单
- 数据量巨大,因此需要在上百台计算机上进行分布式计算
这就提出了几项挑战:
- 如何并行计算
- 如何分发数据
- 如何进行容错
数据处理程序有大量的code都是在处理类似的问题而不是在写业务。
因此,提出了一种新的抽象方式,从lisp和其他的一些functional language中借用了两个概念使得工程师可以专注于业务逻辑,而隐藏上述的非功能性需求,这就是Map-Reduce。
实现
执行逻辑
- 输入数据分片成M份分给M个Map worker
- Master作为数据管道做scheduling
- Map任务调用用户Map方法做完存储在本地内存中
- 周期性dump to disk再分成R份准备分给R个Reduce worker
- Reduce worker在Master给分配任务后通过RPC读数据,读完disk sort
- Reduce worker在sorted数据上调用用户Reduce方法进行iterate,做增量式计算
- 算完收工,回到用户代码逻辑
可以看出来Partition的方法很重要,最好能避免不必要的shuffling
数据结构
- 对每个map/reduce任务存储状态(idle/inprogress/completed)
- 对每个完成的map任务存储中间文件的位置和大小,并推送给reduce worker
容错
- Master基本上不会错,当然也可以采用周期性存储状态的方式在出错后恢复
- Master周期性ping worker,如果ping不通就把worker当做失效worker
- 对于map任务无论完成与否全部重制并重新调度 - 这是因为map的结果会存在worker本地,因此worker失效,结果也访问不到
- 对于reduce任务,只重新调度没完成的 - 这是因为reduce的结果会放在一个全局可访问的文件系统中
- map任务被重新调度之后所有的reduce worker都会收到通知,因此会从新的位置拉取数据
备份任务
最后几个正在执行的任务可能严重拖慢MR节奏,可能的原因是会有“掉队者“,出于各种原因使得他们的执行速度奇慢,这时候master会对最后几个正在执行的任务开启备份任务,这会将资源使用提高几个百分点但是可以海量减少执行时间
改进
Partition
因为输出文件是分布式的,因此可能通过用户提供的partition函数将相关key的结果聚合在一起
本地调试
可以本地运行
性能
在两个大集群中进行测试,一个做sort一个做模式匹配
grep
10^10的100byte文件找pattern,1.7k机器跑了150秒其中分发程序用了一分钟,真实诚
sort
10^10个100byte文件
总结
底气十足