【问题标题】:How to use watchfornewfiles in Dataflow with GCS source bucket?如何在带有 GCS 源存储桶的 Dataflow 中使用 watchfornewfiles?
【发布时间】:2018-06-04 20:54:52
【问题描述】:

参考项目:Watching for new files matching a filepattern in Apache Beam

您可以将它用于简单的用例吗?我的用例是我让用户将数据上传到 Cloud Storage -> Pipeline (Process csv to json) -> Big Query。我知道云存储是有界集合,因此它代表批处理数据流。

我想做的是让管道以流模式运行,一旦文件上传到 Cloud Storage,它将通过管道进行处理。 watchfornewfiles 可以做到这一点吗?

我的代码如下:

p.apply(TextIO.read().from("<bucketname>")         
    .watchForNewFiles(
        // Check for new files every 30 seconds         
        Duration.standardSeconds(30),                      
        // Never stop checking for new files
        Watch.Growth.<String>never()));

没有任何内容被转发到 Big Query,但管道显示它正在流式传输。

【问题讨论】:

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


    【解决方案1】:

    您可以在此处使用 Google Cloud Storage 触发器: https://cloud.google.com/functions/docs/calling/storage#functions-calling-storage-python

    这些触发器使用类似于 Cloud Pub/Sub 的 Cloud Functions,如果对象是:创建/删除/存档/或元数据更改,则会在对象上触发。

    这些事件是使用来自 Cloud Storage 的 Pub/Sub 通知发送的,但请注意不要在同一个存储桶上设置多个函数,因为存在一些通知限制。

    此外,文档末尾还有一个示例实现的链接。

    【讨论】:

    • 将链接中的相关内容添加到答案正文中。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-09-17
    • 2020-06-30
    • 2021-10-24
    • 1970-01-01
    • 1970-01-01
    • 2019-12-25
    • 2021-02-21
    相关资源
    最近更新 更多