【问题标题】:How to sort an arbitrarily large set of data using Hadoop?如何使用 Hadoop 对任意大的数据集进行排序?
【发布时间】:2015-11-06 23:24:28
【问题描述】:

我的问题与Sorting large data using MapReduce/Hadoop 这篇帖子有关。 我对任意集合进行排序的想法是:

  1. 我们有一个包含记录的大文件,比如 10^9 条记录。
  2. 文件在 M 个映射器中拆分。每个映射器使用 QuickSort 对一个大小的拆分(例如 10000 条记录)进行排序,并输出该排序后的子序列。输出键的范围在 1 和 R 之间,其中 R 是 reducer 任务的数量(假设 R = 4)。该值是排序后的子序列。
  3. 每个 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 基准测试。

标签: sorting hadoop mapreduce


【解决方案1】:

Hadoop 排序是使用partitioner 完成的。例如,请参阅 source codeterasort benchmark

【讨论】:

    猜你喜欢
    • 2011-04-07
    • 2016-06-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-11-10
    • 1970-01-01
    • 2014-06-01
    相关资源
    最近更新 更多