【问题标题】:Perform action after Dataflow pipeline has processed all data在 Dataflow 管道处理完所有数据后执行操作
【发布时间】:2017-06-01 19:00:12
【问题描述】:

是否可以在批处理 Dataflow 作业处理完所有数据后执行操作?具体来说,我想将管道刚刚处理的文本文件移动到不同的 GCS 存储桶。我不确定将其放置在我的管道中的哪个位置,以确保它在数据处理完成后执行一次。

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    我不明白您为什么需要执行此后期管道执行。您可以使用侧输出将文件写入多个存储桶,并在管道完成后保存自己的副本。

    如果这对您不起作用(无论出于何种原因),那么您可以简单地在blocking execution 模式下运行您的管道,即使用pipeline.run().waitUntilFinish(),然后只需编写其余代码(执行复制)之后。

    [..]
    // do some stuff before the pipeline runs
    Pipeline pipeline = ...
    pipeline.run().waitUntilFinish();
    // do something after the pipeline finishes here
    [..]
    

    【讨论】:

    • 谢谢。使用BlockingDataflowPipelineRunner 运行作业就可以了。 waitUntilFinish() 在 1.x Java API 中似乎不可用。
    • 正确,不是。您在 1.x 中使用 Blocking runner 并等待/轮询
    • 谢谢波莉。我正在为这些要求创建管道,这完成了我的工作!
    • 如果管道在流模式下执行,这仍然有效吗?这里的用例是,如果管道被窗口化,例如 1 小时间隔,并且我们希望在每个窗口完成后执行一些逻辑。
    【解决方案2】:

    我从阅读apache beam的PassThroughThenCleanup.java的源代码中得到的一个小技巧。

    在您的阅读器之后,创建一个“组合”整个集合的侧输入(在源代码中,它是View.asIterable() PTransform)并将其输出连接到DoFn。这个DoFn 只有在读者读完所有元素后才会被调用。

    附:代码按字面意思命名操作,cleanupSignalView 我发现它非常聪明

    请注意,您可以使用Combine.globally() (java) 或beam.CombineGlobally() (python) 实现相同的效果。欲了解更多信息,请查看第 4.2.4.3 节 here

    【讨论】:

      【解决方案3】:

      我认为这里有两个选项可以帮助您:

      1) 使用TextIO 写入您想要的存储桶或文件夹,指定确切的 GCS 路径(例如 gs://sandbox/other-bucket)

      2) 将Object Change NotificationsCloud Functions 结合使用。您可以在 here 和 JS 中的 GCS 开发工具包 here 中找到很好的入门指南。您在此选项中所做的基本上是在某个存储桶中掉落某物时设置触发器,然后使用您自己编写的 Cloud Function 将其移动到另一个存储桶。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2013-01-03
        相关资源
        最近更新 更多