【问题标题】:Executing a pipeline only after another one finishes on google dataflow只有在另一个管道完成谷歌数据流后才执行管道
【发布时间】:2018-03-09 10:21:56
【问题描述】:

我想在 google 数据流上运行一个依赖于另一个管道输出的管道。现在我只是在本地使用 DirectRunner 运行两个管道:

with beam.Pipeline(options=pipeline_options) as p:
    (p
     | beam.io.ReadFromText(known_args.input)
     | SomeTransform()
     | beam.io.WriteToText('temp'))

with beam.Pipeline(options=pipeline_options) as p:
    (p
     | beam.io.ReadFromText('temp*')
     | AnotherTransform()
     | beam.io.WriteToText(known_args.output))

我的问题如下:

  • DataflowRunner 是否保证仅在第一个管道完成后才启动第二个?
  • 是否有首选的方法来依次运行两条管道?
  • 还有没有推荐的方法将这些管道分成不同的文件以便更好地测试它们?

【问题讨论】:

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


    【解决方案1】:

    DataflowRunner 是否保证第二个仅在第一个管道完成后才启动?

    不,Dataflow 只是执行一个管道。它没有管理依赖管道执行的功能。

    更新:为了澄清,Apache Beam 确实提供了一种等待管道完成执行的机制。请参阅PipelineResult 类的waitUntilFinish() 方法。参考:PipelineResult.waitUntilFinish()

    有没有首选的方法来依次运行两条管道?

    考虑使用 Apache Airflow 之类的工具来管理相关管道。您甚至可以实现一个简单的 bash 脚本来在另一个管道完成后部署一个管道。

    还有没有推荐的方法将这些管道分成不同的文件以便更好地测试它们?

    是的,单独的文件。这只是很好的代码组织,不一定更适合测试。

    【讨论】:

    • 您的第一点并不完全正确。确实,Dataflow 本身并不管理管道依赖项,但它确实提供了一个 API 来等待管道完成,并且 Pipeline 对象上的“with”语句使用该 API。
    • @jkff,是的,没错。 Apache Beam Java SDK 中的PipelineResult 类有一个方法waitUntilFinish()
    • python呢?
    猜你喜欢
    • 1970-01-01
    • 2017-05-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-02-28
    • 1970-01-01
    • 2022-08-17
    • 1970-01-01
    相关资源
    最近更新 更多