【问题标题】:Dataflow Distinct operation not scaling数据流不同的操作不缩放
【发布时间】:2018-07-14 05:31:07
【问题描述】:

我有一个带有“最终”阶段的线性管道,每秒输出大约 200k 个元素(短字符串)。

但是,当我在该阶段 (myPCollection.apply(Distinct.<String>create());) 之后添加 Distinct 操作时,Distinct 之前的阶段速度下降到每秒处理的大约 80k 个元素。

但是,我正在处理一个没有最大工作人员数量的有界集合,因此我希望 Dataflow 能够自动增加工作人员的数量以匹配工作负载。不仅不会发生这种情况,而且当我手动启动具有许多工作人员(20 多名)的管道时,它会自动缩减为几个工作人员。

如何使 Dataflow 升级工作器池,以便此 Distinct 操作不会显着降低管道的处理速率?

【问题讨论】:

    标签: performance google-cloud-dataflow apache-beam autoscaling


    【解决方案1】:

    看看the implementation of Distinct可能会很有趣。

    如您所见,它首先对元素进行分组,然后再选取第一个元素。我已经提交了a bug 以改善这种行为。

    在当前实现中所有元素首先被分组,这需要将它们写入持久存储,然后再被拾取。如果您有任何元素多次出现(即一个热键),您将遇到可以写出多少数据的瓶颈。

    作为一个技巧,您可以在写出元素之前添加一个 DoFn 去重复元素。像这样的:

    class MapperDedupFn extends DoFn<String, String> {
        Set<String> seenElements;
        MapperDedupFn() {
          seenElements = new HashSet<>();
        }
    
        @ProcessElement
        public void processElement(@Element String element, OutputReceiver<String> receiver) {
          if (seenElements.contains(element)) return;
    
          seenElements.add(element)
          receiver.output(word);
        }
      }
    }
    

    你应该可以在Distinct函数之前坚持这个,希望有更好的性能。

    【讨论】:

    • 非常感谢您的回答。然而,在稍微重构了我的代码并完全删除了Distinct 操作之后,我意识到我的管道仍然无法扩展,即使某些阶段积压了数百万个项目!我的管道仅由几个线性方式的 DoFns 组成。想想 TextIO.read(快)——一些处理 DoFns(快)——匹配的 DoFn(慢)——TextIO.write(快)。即使是 DoFn,管道也无法扩展以弥补缓慢的 DoFn。您是否有任何见解可以解释这种行为?非常感谢。
    • 很高兴为您提供帮助。你所说的“匹配”是什么意思?随时将您的问题发送至user@beam.apache.org,或开始一个我可以在 SO 上回答的新问题(在此处链接以便我找到它)
    • 经过一些研究,我认为问题来自一对多阶段,并已发起new SO question。谢谢
    猜你喜欢
    • 1970-01-01
    • 2020-12-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多