【问题标题】:How to read from PubSub in DataFlow using batches如何使用批处理从 DataFlow 中的 PubSub 读取
【发布时间】:2018-05-13 10:03:47
【问题描述】:

在 Pubsub 源的 SDK 1.9.1 中,有可用的 PubsubIO.Read.maxReadTimePubsubIO.Read.maxNumRecords 方法。这些方法允许从 pubsub 消息创建有界集合,可以在批处理模式下启动 Dataflow 管道。

使用 Dataflow SDK 2.1 可以实现多少类似的事情?如何使用批处理模式从 Dataflow 管道中的 Pubsub 读取?

【问题讨论】:

  • 你检查过这个吗? beam.apache.org/documentation/sdks/javadoc/2.1.0/index.html?org/…我最近使用了更多的 Scio,但我在“纯”Beam 中记不太清了,但它似乎与您正在寻找的相似。但是似乎它必须放在 PubsubIO.Read 之后
  • 确实看起来非常相似,但是如何将其应用于管道?使用它的唯一方法是访问深埋在SDK中的源代码,并且在与apply一起使用后会丢失。提供 PubsubIO.Read 的开发人员可以使用它,但是使用 PubsubIO.Read API 的开发人员如何使用它?

标签: google-cloud-platform google-cloud-dataflow google-cloud-pubsub gcp


【解决方案1】:

很遗憾,我在新版本的 SDK 中没有看到任何支持。我所做的是实现一个 DoFn,它从 PubSub 读取 ma​​xReadTimema​​xNumRecords 并返回消息。

这就是他们在以前版本的 SDK 上所做的。您可以查看PubsubReader 类。

你必须这样称呼它:

 pipeline.begin()
            .apply(Create.of((Void) null)).setCoder(VoidCoder.of())
            .apply(ParDo. of(new MyPubsubReader(maxNumRecords, maxReadTime));
            .setCoder(coder);

【讨论】:

    【解决方案2】:

    您不应尝试在批处理上下文中使用 PubsubReader。相反,您应该使用提供的流式 PubsubIO,并按照here 的描述设置窗口策略。您可以使用“其他复合触发器”部分(复制如下)中描述的复合触发器来获得所需的行为。

    Repeatedly.forever(AfterFirst.of(
          AfterPane.elementCountAtLeast(100),
          AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(1))))
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-07-27
      • 1970-01-01
      • 2018-08-10
      • 2020-10-16
      • 2011-12-04
      • 2021-09-05
      • 1970-01-01
      • 2021-12-15
      相关资源
      最近更新 更多