【问题标题】:Why is groupBy bottlenecking my pipeline?为什么 groupBy 会阻碍我的管道?
【发布时间】: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


【解决方案1】:

总结jkff等人的回答:

管道似乎受到单个非常大的键的限制。您可以使用常规 Java 日志记录并查看工作日志(例如,在 processElement() 中测量 DoFn 的处理时间,如果超过阈值则记录它),但不幸的是,我们还没有提供更高级别的工具来调试“热键”问题。

您还可以打开autoscaling,以便该服务至少可以关闭未使用的工作人员,这样您就不会为他们产生费用。

【讨论】:

  • 很抱歉问了一个老问题,找不到更好的地方。 Window+GroupByKey 的一对变换不是使用“Window+Key”作为分组操作的键吗?例如,如果输入集合有 800k 个时间序列元素,每秒一个元素,并且它被窗口化,比如 1 分钟的窗口,然后按 10 个键分组 - 它是否会生成 10 个组,然后将它们拆分为触发单独的窗口,还是一次分组为 800'000/60 / 10 ~= 1300 个组?因为上面的 cmets(大约一个大键)导致想到前者,但这将是一个令人惊讶的行为
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-04-21
  • 1970-01-01
  • 2023-02-11
相关资源
最近更新 更多