【发布时间】:2015-11-06 23:24:28
【问题描述】:
我的问题与Sorting large data using MapReduce/Hadoop 这篇帖子有关。 我对任意集合进行排序的想法是:
- 我们有一个包含记录的大文件,比如 10^9 条记录。
- 文件在 M 个映射器中拆分。每个映射器使用 QuickSort 对一个大小的拆分(例如 10000 条记录)进行排序,并输出该排序后的子序列。输出键的范围在 1 和 R 之间,其中 R 是 reducer 任务的数量(假设 R = 4)。该值是排序后的子序列。
- 每个 Reducer 读取 K 个子序列并将它们合并(迭代地从子序列中取出最小元素,直到子序列为空)。输出被写入文件。
然后进行如下处理:
为了利用数据的局部性,可以安排新的 Reducer 任务来合并前一个 reducer 任务生成的多个输出文件。因此,例如,如果 K=5,第一个 reducer 任务将生成大小为 50000 的文件,而新的 reducer 任务将处理 5 个文件,每个文件包含 50000 个排序记录。新的 Reducer 作业将被安排到只剩下一个文件,在这种情况下大小为 250.000.000(因为 R=4)。最后,将在另一台机器上安排一个新的 Reducer 作业,将文件合并为一个 10^9 文件
我的问题:是否可以在 Hadoop 中安排 Reducer 作业的执行,以便它们合并某个目录中的文件,直到只剩下 1 个文件?如果是,怎么做?
另一种情况是在每个合并步骤之后安排 MapReduce 作业,例如,大小为 50000 的文件将通过在其他机器上运行的 reduce 任务并行合并,然后在其他机器上运行大小为 250.000 的文件,等等。但这会产生大量的网络流量。无论如何,这个问题对于这种情况也仍然有效 - 如何链接多个 MapReduce 作业,以便在仅输出 1 个结果文件后链接停止?
【问题讨论】:
-
没有任何开销,它可以为您完成所有工作,而且确实是最佳选择。他们赢得了 terasort 基准测试。