【问题标题】:apache beam global combine shuffleapache Beam 全局组合洗牌
【发布时间】:2018-08-28 09:52:04
【问题描述】:

在我的 apache 梁和数据流管道中,我做了一些需要全局合并操作的转换,例如 min 、 max 、自定义全局合并函数。 pcollection 中要处理的项目数量为 2-40 亿。

我读到大多数组合操作都是建立在 groupBykey 之上的,这会导致 shuffle ,我相信这会使我当前的管道变慢,或者从 UI 中观察到,这是全局组合操作中最高的挂壁时间。我查看了代码, groupByKey 尝试向所有元素添加一个静态 void 键,然后执行 groupby ,这是否意味着我们正在改组数据(特别是当我们只有一个键时)? 有没有办法有效地做到这一点

我自己理解的另一个问题:beam/dataflow 文档说键的所有元素都由单个工作线程/线程处理。以在整数的 pcollection 中查找最大值为例,此全局操作是完全可并行化的,其中我的组合器/累加器在数据的部分/子集上工作以找到最大值,然后在树中合并部分结果(合并两个最大值以获得最大值)类似结构,叶子的结果可以合并得到父节点,每个节点基本上可以分布式评估。那么究竟是什么操作强制一个键必须由一个工作线程/线程处理。似乎任何具有可交换和关联组合器的全局操作都可以轻松并行化。全局组合的哪一部分需要通过单个工作线程?

【问题讨论】:

  • 我已经观察到这种行为,我正在努力解决这个问题。当您创建一个 Combine.globally() 时,只有方法 createAccumulator、addInput 和 extractOutput 被调用。这意味着 threre 是一个 shuffle 并且不需要 mergeAccumulators。我试图找出寻找代码的原因,但 apache 梁代码很糟糕。对我来说,在以后的部分合并中对每个捆绑包进行部分合并并最终收集所有内容是有意义的。但这并不是真正发生的事情,所以医生在撒谎。

标签: google-cloud-platform google-cloud-dataflow shuffle apache-beam apache-beam-io


【解决方案1】:

组合器将在随机播放之前被提升(这意味着我们在传递给随机播放之前会进行一些组合)。这里有一点信息:https://cloud.google.com/blog/big-data/2016/02/writing-dataflow-pipelines-with-scalability-in-mind,搜索combiner。

Dataflow 将为每个元素分配一个不同的键,因此您最终不会得到所有相同的键(因此没有并行性)。如果全部分配给一个key,那么只有一个worker可以处理,而且会很慢。

【讨论】:

  • 所以基本上任何全局合并操作都会太慢?
  • 即使找到最小值和最大值也会太慢?通常不会发生这种情况
  • 你能解释一下每个元素的不同键是什么意思吗,源代码显示它为每个元素分配了相同的 void 键?我对你所说的“所以你不会得到所有相同的键(因此没有并行性)”感到非常困惑。对我来说似乎是错误的,更多的键意味着更多的洗牌和更多的并行性?
  • 我不知道代码显示了什么,但结果不是每个元素上的键都相同。稍后在管道中可能会替换无效键。如果每个元素都有相同的键,您将不会获得并行性,它将是单线程的。全局合并将有一个洗牌,但不应该太慢,因为一些工作是在洗牌之前完成的。你是什​​么意思'通常不会发生'。
  • 好的阅读 beam/dataflow source code ,我们添加 None (python) 和 Void ,每个元素的键。这与我在查看自动缩放的工人数量时观察到的一致。在我的全球联合作业中,工人人数减少到 1 人。另一个问题是,为什么你需要一个 shuffle 来进行 global combine ,为什么你需要一个 shuffle 数据,当你只有一个 key 时,它背后的原理是什么?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-11-10
  • 1970-01-01
  • 2018-08-06
  • 1970-01-01
  • 2021-10-20
相关资源
最近更新 更多