【问题标题】:Parallel and bulk process datastore elements in google cloud dataflow谷歌云数据流中的并行和批量处理数据存储元素
【发布时间】:2017-09-13 13:27:14
【问题描述】:

问题:

我的数据存储项目中有超过 200 万个用户数据的列表。我想向所有用户发送每周通讯。邮件 API 每个 API 调用最多接受 50 个电子邮件地址。

以前的解决方案:

使用应用引擎后端和简单的数据存储查询一次性处理所有记录。但是发生的情况是,有时我会收到内存溢出严重错误日志,并且该过程会重新开始。因此,一些用户会多次收到同一封电子邮件。所以我转向数据流。

当前解决方案:

我使用 FlatMap 函数将每个电子邮件 ID 发送到一个函数,然后将电子邮件单独发送给每个用户。

def process_datastore(project, pipeline_options):
    p = beam.Pipeline(options=pipeline_options)
    query = make_query()
    entities = (p | 'read from datastore' >> ReadFromDatastore(project, query))
    entities | beam.FlatMap(lambda entity: sendMail([entity.properties.get('emailID', "")]))
    return p.run()

通过云数据流,我确保每个用户只收到一次邮件,并且没有人错过。没有内存错误。

但是这个当前进程需要 7 个小时才能完成运行。我试图用 ParDo 替换 FlatMap,假设 ParDo 将并行化该过程。但即使这样也需要同样的时间。

问题:

  1. 如何将50个email ids分组,以有效利用邮件API调用?

  2. 如何并行处理,使耗时少于一小时?

【问题讨论】:

标签: google-app-engine google-cloud-dataflow apache-beam


【解决方案1】:

您可以使用查询游标将用户分成 50 个批次,并在推送队列或deferred 任务中进行实际的批处理(电子邮件发送)。这将是一个仅限 GAE 的解决方案,没有云数据流,恕我直言要简单得多。

您可以在Google appengine: Task queue performance 中找到此类处理的示例(也将答案考虑在内)。该解决方案使用deferred 库,但使用推送队列任务几乎是微不足道的。

答案涉及并行性方面,您可能希望限制并行性以降低成本。

您还可以将批处理本身拆分到任务中以获得无限可扩展的解决方案(任意数量的接收者,而不会遇到内存或超出期限的故障),任务重新排队以从中断的地方继续工作.

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-11-23
    • 1970-01-01
    • 2020-08-04
    • 2015-07-17
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多