【问题标题】:How to propagate PubSub metadata with Apache Beam?如何使用 Apache Beam 传播 PubSub 元数据?
【发布时间】:2020-01-10 09:28:11
【问题描述】:

上下文:我有一个监听 pub sub 的管道,发送到 pubsub 的消息是由来自谷歌云存储的对象更改通知发布的。管道使用 XmlIO 拆分文件来处理文件,到目前为止一切顺利。

问题是:在 pubsub 消息中(以及存储在谷歌云存储中的对象中)我有一些元数据,我想与来自 XmlIO 的数据合并以组成管道将处理的元素,如何我能做到吗?

【问题讨论】:

  • 您是否在管道中使用任何窗口/触发?
  • 我没有使用窗口/触发

标签: google-cloud-storage google-cloud-dataflow apache-beam


【解决方案1】:

您可以创建一个自定义窗口和 windowfn 来存储来自 pubsub 消息的元数据,以便以后使用这些元数据来丰富各个记录。

您的管道将如下所示:

ReadFromPubsub -> Window.into(CopyMetadataToCustomWindowFn) -> ParDo(ExtractFilenameFromPubsubMessage) -> XmlIO -> ParDo(EnrichRecordsWithWindowMetadata) -> Window.into(FixedWindows.of(...))

首先,您需要创建一个IntervalWindow 的子类来存储您需要的元数据。之后,创建WindowFn 的子类,其中#assignWindows(...) 将元数据从pubsub 消息复制到您创建的IntervalWindow 子类中。使用 Window.into(...) 转换应用您的新 windowfn。现在,通过 XmlIO 转换的每条记录都将位于包含元数据的自定义 windowfn 中。

对于第二步,您需要从 pubsub 消息中提取相关文件名,以作为输入传递给 XmlIO 转换。

对于第三步,您希望从位于 XmlIO 之后的 ParDo/DoFn 中的窗口中提取自定义元数据。 XmlIO 中的记录将保留通过它传递的窗口信息(请注意,并非所有转换都这样做,但几乎所有转换都这样做)。您可以声明您的DoFn needs the window to be passed to your @ProcessElement,例如:

class EnrichRecordsWithWindowMetadata extends DoFn<...> {
  @ProcessElement
  public void processElement(@Element XmlRecord xmlRecord, MyCustomMetadataWindow metadataWindow) {
    ... enrich record with metadata on window ...
  }
}

最后,恢复到标准 windowfns 之一是个好主意,例如 FixedWindows,因为窗口上的元数据不再相关。

【讨论】:

    【解决方案2】:

    您可以直接使用来自 Google Cloud Storage 的 pub/sub 通知,而不是在中间引入 OCN。

    Google 也建议使用 pub/sub。如果您收到发布/订阅通知,您可以在其中获取消息属性。

    data = request.get_json()
    
    object_id = data['message']['attributes']['objectGeneration']
    bucket_name = data['message']['attributes']['bucketId']
    object_name = data['message']['attributes']['objectId']
    

    【讨论】:

    • 这个我知道,问题是如何在数据流管道中将此元数据与 XmlIO 的结果合并。
    猜你喜欢
    • 2020-06-28
    • 1970-01-01
    • 2019-06-04
    • 1970-01-01
    • 1970-01-01
    • 2022-12-24
    • 1970-01-01
    • 2021-07-24
    • 2018-12-30
    相关资源
    最近更新 更多