【发布时间】: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