【发布时间】:2016-09-22 22:09:18
【问题描述】:
我有一个用 python apache-beam 编写的管道。它将 800,000 个时间戳数据窗口化为每 1 秒重叠的 2 秒窗口。我的元素可能有不同的键。
当它执行 groupBy 时,需要 3 个小时才能完成。我使用 10 个工作人员部署在云数据流中。当我增加工人数量时,处理速度并没有显着提高。为什么这种转变会成为我的管道的瓶颈?
【问题讨论】:
-
我不熟悉这个问题的具体标签,但是为了执行“groupby”操作(在非索引字段上),那么分组必须对整个数据集在一个地方。如果没有该操作,您可能拥有易于并行化的数据。是这个词吗?
-
是否有特定的作业 ID 证明了问题?正如 Kenny Ostrom 所提到的,gorupby 不仅需要在一个地方对数据进行排序,还需要在工作人员之间传递该数据,以便处理单个键的工作人员拥有所有关联的值。元素有多大(这应该在每个步骤的数据流 UI 中显示)?
-
除非元素很大,否则对 800,000 个元素本身进行分组应该非常快。很可能还有其他东西阻碍了管道 - 您是否在为每个元素或每个键做一些非常昂贵的事情?有多少种不同的键? (单个键是按顺序处理的,因此如果键很少,那么无论您指定多少工作人员,这都会限制可实现的最大并行度)确实,如果没有要查看的作业 ID,就很难判断发生了什么。
-
你用什么做key?密钥是如何计算的?您的所有(或大部分)数据是否可能都在同一个密钥上?与键关联的所有元素都需要在单个工作器上按顺序处理。
-
您可以使用常规 Java 日志记录并查看工作日志(例如,在 processElement() 中测量您的 DoFn 的处理时间,如果超过阈值则将其记录下来),但遗憾的是我们还没有提供更高的 -用于调试“热键”问题的级别工具。我查看了这些管道,确实,它们都被一个非常大的密钥有效地限制了。我还建议打开自动缩放,以便服务至少可以关闭未使用的工作人员,这样您就不会为他们产生费用。
标签: python google-cloud-dataflow dataflow apache-beam