【问题标题】:Google Dataflow: Read unbound PCollection from Google Cloud StorageGoogle Dataflow:从 Google Cloud Storage 读取未绑定的 PCollection
【发布时间】:2020-02-27 11:47:18
【问题描述】:

我有一个将 JSON 消息从 PubSub(未绑定 PCollection)流式传输到 Google Cloud Storage 的管道。每个文件应包含多个 JSON 对象,每行一个。

我想创建另一个管道,该管道应该从这个 GCS 存储桶中读取所有 JSON 对象,以进行进一步的流处理。最重要的是,第二个管道应该作为流而不是批处理。意味着我希望它“监听”存储桶并处理写入其中的每个 JSON 对象。 未绑定 PCollection。

有没有办法实现这种行为?

谢谢

【问题讨论】:

    标签: java google-cloud-storage google-cloud-dataflow stream-processing


    【解决方案1】:

    流式处理仅适用于 PubSub 数据源。不过不用担心,您可以实现自己的管道。

    【讨论】:

    • 我很想知道您在使用这种方法时所体验到的性能。我正在采取类似的方法,而 beam.io.ReadAllFromText() 步骤似乎是一个真正的瓶颈。实现令人满意的吞吐量所需的工人数量是不可行的。
    【解决方案2】:

    另一个用户提供了一个很好的答案,告诉你如何做你想做的事,但是,如果我正确理解你的问题,我想我可以推荐一种更简洁的方法。

    假设以下情况属实,您希望:

    • 在管道开始时接受 Pub/Sub 消息
    • 通过将消息窗口化到消息的大窗口中来处理消息
    • 将每个消息窗口作为文件写入 GCS 存储桶
    • 除了上面描述的窗口化之外,以另一种方式处理每一行

    然后,您可以创建一个在“接受 Pub/Sub 消息”步骤之后简单地分叉的管道。 Dataflow 本身就很好地支持了这一点。您将保存对在管道开始处使用 Pub/Sub 接收器时返回的 PCollection 对象的引用。然后,您可以将多个DoFn 实现链等应用到这个引用中。您将能够像现在一样通过写入 GCS 来进行窗口化,并以您喜欢的任何方式处理每条单独的消息。

    它可能看起来像这样:

    Pipeline pipeline = Pipeline.create(options);
    
    PCollection<String> messages = pipeline.apply("Read from Pub/Sub", PubsubIO.readStrings().fromTopic("my_topic_name));
    
    // Current pipeline fork for windowing into GCS
    messages.apply("Handle for GCS", ...);
    
    // New fork for more handling
    messages.apply("More stuff", ...);
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-08-01
      • 2019-11-17
      • 1970-01-01
      • 1970-01-01
      • 2018-07-22
      • 1970-01-01
      相关资源
      最近更新 更多