【问题标题】:Early results from GroupByKey transformGroupByKey 转换的早期结果
【发布时间】:2018-02-20 13:57:15
【问题描述】:

我怎样才能让 GroupByKey 触发早期结果,而不是等待所有数据到达(在我的情况下这是一个相当长的时间)。我尝试将我的输入 PCollection 拆分为具有早期触发器的窗口,但它只是行不通。在给出结果之前,它仍然等待所有数据到达。

PCollection<List<String>> input = ...
PCollection<KV<Integer,List<String>>> keyedInput = input.apply(ParDo.of(new AddArbitraryKey()))
keyedInput.apply(Window<KV<Integer,List<String>>>into(
          FixedWindows.of(Duration.standardSeconds(1)))
         .triggering(Repeatedly.forever(AfterWatermark.pastEndOfWindow()))
         .withAllowedLateness(Duration.ZERO).discardingFiredPanes())
 .apply(GroupByKey.<Integer,List<String>>create())
       .apply(ParDo.of(new RemoveArbitraryKey()))
       .apply(ParDo.of(new FurtherProcessing())

我这样做是为了防止 fusing 。 AddArbitraryKey 转换使用时间戳输出其元素。但是,GroupByKey 会保留所有内容,直到所有数据到达(对于所有窗口)。有人可以告诉我如何让它尽早触发。谢谢你 。

【问题讨论】:

  • 我刚刚遇到了同样的问题,一个简单的固定窗口没有帮助,一个组合每个键等到有界集合结束。你找到解决办法了吗?

标签: google-cloud-dataflow apache-beam


【解决方案1】:

你可以安装一个触发器像

Repeatedly
  .forever(AfterProcessingTime
    .pastFirstElementInPane()
    .plusDuration(Duration.standardMinutes(1))
  .orFinally(AfterWatermark.pastEndOfWindow())
  .discardingFiredPanes()

或者

AfterWatermark.pastEndOfWindow()
  .withEarlyFirings(
    AfterProcessingTime
      .pastFirstElementInPane()
      .plusDuration(Duration.standardMinutes(1))

【讨论】:

  • 这对我不起作用。 GroupBy 转换仅在所有元素到达后​​才分派结果。我将窗口和触发器应用于 PCollection,就在 groupby 转换之前。那是对的吗 ?我的意图是,一旦一个时间窗口内的所有元素都到达,GroupBy 应该为该窗口发出其结果。我有错误的想法吗?或者您能告诉我如何实现这一目标吗?
  • 你试过.plusDuration(Duration.ZERO)吗?另外,你的数据源是什么?是批量数据源吗?
  • 是的,我试过 Duration.ZERO 。是的,它是批处理数据源,而不是流式传输。所以,PCollection 是有界的。在这种情况下会起作用吗?
  • 批处理管道通常不需要窗口化。 GBK 将等待所有数据,直到发出输出。如果您仍需要流式处理行为(窗口和触发),则需要使用 outputWithTimestamp 为每个输入元素注入事件时间戳(不是处理时间)。
  • 我在 ParDo 中使用了 outputWithTimestamp:processContext.outputWithTimestamp(id,new Instant()); .那是对的吗 ?即使在有界 PCollection 上使用 Windows 和触发器,GBK 也会等待所有数据
【解决方案2】:

为了防止融合,最好使用转换Reshuffle.viaRandomKey(),它性能更好,并确保不会引入任何额外的触发延迟。

【讨论】:

  • Reshuffle 中的 GroupBy 转换成为瓶颈,即它一直等到所有元素都到达。在应用 Reshuffle.viaRandomKey 之前,我对元素加了时间戳并将它们应用到一个窗口(如我的问题中的代码 sn-p 所示)。我的意图是,一旦时间窗口内的所有元素到达,GroupBy 应该为该窗口发出其结果。我有错误的想法吗?
  • 我正在使用的 PCollection 是一个有界 PCollection(批处理管道的)。这有什么不同吗?
  • 批处理管道针对吞吐量而不是延迟进行优化,并且没有水印跟踪(部分原因是像文件这样的典型批处理数据源没有明显的时间戳排序并且无法提供任何有用的水印估计) - 所以GroupByKey 有效地缓冲所有数据并在所有数据到达时触发所有窗口。这是否显示为您的管道的性能问题?
  • 是的,在我的场景中:“源”是数据流程序(从 ParDo 中)命中以获取记录的 Web 服务。ParDo 内部有一个循环,在其中它调用 Web 服务. Web 服务在其末尾查询数据库,并使用游标以顺序方式但分批返回数据。因此,这个特殊的 ParDo 在单个工作人员上运行。这具有高扇出,因为在循环内为批处理中返回的每个记录调用 output()。处理这些记录的下一个 ParDo 不会根据第一个 ParDo 输出的记录数进行缩放。融合?
  • 我怀疑可能正在发生融合,这会阻止下一个 ParDo 扩展,因此要打破它,想要分组/取消分组。但是如果 GBK 等待所有数据到达,它会阻止阻止已经到达的记录被传递到下一个 ParDo。 (网络服务需要很长时间才能完成所有批次,在此之前没有任何记录被进一步处理)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-08-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2010-11-17
相关资源
最近更新 更多