【问题标题】:Processing a large amount of data in parallel并行处理大量数据
【发布时间】:2012-12-22 20:30:32
【问题描述】:

我是一名 Python 开发人员,拥有很好的 RDBMS 经验。我需要处理相当大量的数据(大约 500GB)。数据位于 s3 存储桶中的大约 1200 个 csv 文件中。我用 Python 编写了一个脚本,可以在服务器上运行它。但是,速度太慢了。根据目前的速度和数据量,通过所有文件大约需要 50 天(当然,截止日期在此之前是很好的)。

注意:处理是您的基本 ETL 类型的东西 - 没有什么可怕的幻想。我可以轻松地将其泵入 PostgreSQL 中的临时模式,然后在其上运行脚本。但是,再一次,从我最初的测试来看,这会很慢。

注意:一个全新的 PostgreSQL 9.1 数据库将是它的最终目的地。

所以,我正在考虑尝试启动一堆 EC2 实例以尝试分批(并行)运行它们。但是,我以前从未做过这样的事情,所以我一直在四处寻找想法等。

再说一次,我是一名 python 开发人员,所以 Fabric + boto 似乎很有希望。我不时使用过boto,但从未使用过Fabric。

我从阅读/研究中知道,这对于 Hadoop 来说可能是一项很棒的工作,但我不知道,也负担不起雇用它完成的费用,而且时间线不允许学习曲线或雇用某人.我也不应该,这是一种一次性的交易。所以,我不需要构建一个非常优雅的解决方案。我只需要让它工作,并能够在年底前获得所有数据。

另外,我知道这不是一个简单的 stackoverflow 问题(类似于“如何在 python 中反转列表”)。但是,我希望有人读到这篇文章并“说,我做了类似的事情并使用 XYZ ......这太棒了!”

我想我要问的是有没有人知道我可以用来完成这项任务的任何东西(鉴于我是一名 Python 开发人员并且我不了解 Hadoop 或 Java - 并且有一个紧张的阻止我学习 Hadoop 等新技术或学习新语言的时间表)

感谢阅读。我期待任何建议。

【问题讨论】:

  • fabric+boto 看起来确实是这个任务的好组合。在每个实例上并行化任务也可能是值得的(除非您期望有 1200 个实例,每个文件一个),也许可以使用来自 multiprocessing 模块的 Pool。此外,您解析文件和编辑结果的方式可能会对总时间产生很大影响。你看过numpy吗?
  • 所以没有人会尝试重复可能的建议 - 您能否描述一下您在现有脚本中所做的工作太慢了 - 所以我们知道不要走那条路 :)
  • @JonClements - 似乎是一个公平的要求。基本上,我尝试了两种方法。我尝试将数据放入临时模式并对其进行索引(根据需要)并对它运行查询以“按摩”数据并将其转换为请求的格式。这太慢了,因为我相信索引比 PostgreSQL 缓存大得多。注意:我有一个在 Heroku 上运行的小型 PostgreSQL 实例。 (将在下一条评论中继续)
  • 然后我尝试将我需要的所有数据加载到 python 字典中并通过将它们从 s3 中拉出并加载它们然后以正确的 CSV 格式将它们吐回然后对目标表进行“复制”。这可行,但只有 1 台服务器太慢。但是,我认为如果我采用它并在 100 台服务器上运行它,它将在我的期限内完成这项工作。
  • 请注意。假期我请了一点假,但现在又开始工作了。项目完成后,我将选择一个答案(或发布我所做的)。感谢大家的 cmets/answers。

标签: python fabric boto data-processing


【解决方案1】:

您是否进行了一些性能测量:瓶颈在哪里?是 CPU 限制、IO 限制还是 DB 限制?

当它受 CPU 限制时,您可以尝试使用 pypy 之类的 python JIT。

当它是 IO 绑定时,您需要更多的 HD(并在它们上放置一些条带化 md)。

当它是DB绑定时,你可以尝试先删除所有的索引和键。

上周我将 Openstreetmap 数据库导入到我服务器上的 postgres 实例中。输入数据约为450G。预处理(在这里用 JAVA 完成)刚刚创建了可以使用 postgres 'copy' 命令导入的原始数据文件。导入后生成键和索引。

导入所有原始数据大约需要一天时间 - 然后需要几天时间来构建键和索引。

【讨论】:

    【解决方案2】:

    我经常将 SQS/S3/EC2 组合用于此类批处理工作。在 SQS 中为所有需要执行的工作排队消息(分成一些相当小的块)。启动 N 个配置为开始从 SQS 读取消息、执行工作并将结果放入 S3 的 N 个 EC2 实例,然后,只有这样,才从 SQS 中删除消息。

    你可以将它扩展到疯狂的水平,它对我来说一直很有效。就您而言,我不知道您是将结果存储在 S3 中还是直接转到 PostgreSQL。

    【讨论】:

    • 只是出于好奇,您如何将脚本传送到 EC2 实例?你会让他们从 git repo 中提取吗?或者只是 scp 脚本结束?
    • 我使用了很多技巧。您可以编写基于 Paramiko 的脚本来 scp 文件。您可以使用 cloud-init 并从 S3 中提取脚本。你可以使用织物。您可以使用 CloudFormation 模板。有很多选择。
    • 感谢您的回复。是的;似乎有很多选择。正如我在最初的问题中提到的,我倾向于使用 Fabric,但想知道你在这里做了什么。
    【解决方案3】:

    前段时间我做过类似的事情,我的设置是这样的

    • 一个多核实例(x-large 或更多),可将原始源文件 (xml/csv) 转换为中间格式。您可以在其上并行运行转换器脚本的 (num-of-cores) 副本。由于我的目标是 mongo,所以我使用 json 作为中间格式,在你的情况下它将是 sql。

    • 此实例附加了 N 个卷。一旦一个卷变满,它就会被分离并附加到第二个实例(通过 boto)。

    • 第二个实例运行一个 DBMS 服务器和一个将准备好的 (sql) 数据导入数据库的脚本。我对 postgres 一无所知,但我猜它确实有一个像 mysqlmongoimport 这样的工具。如果是,请使用它进行批量插入,而不是通过 python 脚本进行查询。

    【讨论】:

      【解决方案4】:

      您可能会从 Amazon Elastic Map Reduce 形式的 hadoop 中受益。无需太深入,它可以被视为一种将一些逻辑应用于并行(Map 阶段)中的海量数据的方法。
      还有一种称为 hadoop 流的 hadoop 技术——它可以使用任何语言(如 python)的脚本/可执行文件。
      您会发现另一种有用的 hadoop 技术是 sqoop - 它在 HDFS 和 RDBMS 之间移动数据。

      【讨论】:

      • 感谢您的回答。在我内心深处,我知道 Hadoop 和 Elastic MapReduce 在这里使用是正确的。但是,我只是无法理解它如何与我想要完成的工作一起工作。我的部分问题是,我见过的几乎每个例子都是同样愚蠢的字数计算问题。我的确实更像是一个 ETL(提取、转换、加载)问题。我可以很容易地想象 Map 函数处理大部分转换。但是,转换取决于客户。所以,这不是一个简单的计算(例如(x*y)/2)。
      【解决方案5】:

      您还可以通过StarClusterStarCluster 在 EC2 上非常轻松地使用 ipython 的并行计算,这是一个用于创建和管理托管在 Amazon 的 EC2 上的分布式计算集群的实用程序。

      http://ipython.org/ipython-doc/stable/parallel/parallel_demos.html
      http://star.mit.edu/cluster/docs/0.93.3/index.html
      http://star.mit.edu/cluster/docs/0.93.3/plugins/ipython.html

      【讨论】:

      • 感谢您的帖子!我从来没有听说过这个!我会检查出来的!再次感谢!
      猜你喜欢
      • 1970-01-01
      • 2019-06-18
      • 2017-02-09
      • 2012-10-27
      • 2013-01-10
      • 2011-01-10
      • 2017-09-26
      • 2021-04-08
      相关资源
      最近更新 更多