【发布时间】:2017-05-24 22:59:28
【问题描述】:
我有一个 Python beam.DoFn,它正在将文件上传到互联网。此过程使用 100% 的一个内核约 5 秒,然后继续上传文件 2-3 分钟(并在上传期间使用非常小的一部分 cpu)。
DataFlow 是否足够聪明,可以通过在单独的线程/进程中启动多个 DoFn 来优化这一点?
【问题讨论】:
我有一个 Python beam.DoFn,它正在将文件上传到互联网。此过程使用 100% 的一个内核约 5 秒,然后继续上传文件 2-3 分钟(并在上传期间使用非常小的一部分 cpu)。
DataFlow 是否足够聪明,可以通过在单独的线程/进程中启动多个 DoFn 来优化这一点?
【问题讨论】:
是的,Dataflow 将使用 python 多处理运行多个 DoFn 实例。
但是,请记住,如果您使用 GroupByKey,则 ParDo 将连续处理特定键的元素。虽然您仍然可以在工作人员上实现并行性,因为您一次处理多个键。但是,如果您的所有数据都在一个“热键”上,您可能无法实现良好的并行性。
您是否在批处理管道中使用 TextIO.Write?我相信这些文件是在本地准备好的,然后在处理完您的主要 DoFn 后上传。即在 PCollection 完成之前不会上传文件,并且不会接收更多元素。
我认为它不会在您生成元素时流出文件。
【讨论】: