【发布时间】: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 将并行化该过程。但即使这样也需要同样的时间。
问题:
如何将50个email ids分组,以有效利用邮件API调用?
如何并行处理,使耗时少于一小时?
【问题讨论】:
-
您的管道可能遇到github.com/apache/beam/blob/master/sdks/python/apache_beam/io/… 此处描述的一种情况。在这种情况下,您需要通过打破融合来引入更多并行性cloud.google.com/dataflow/service/…
-
另外:FlatMap 和 ParDo 是等价的——它们都是平行的,但都需要融合(见上文)。一般来说,在调试作业性能时,请提供数据流作业 ID,以便 oncall 工程师查看。
标签: google-app-engine google-cloud-dataflow apache-beam