【发布时间】:2016-11-22 00:22:24
【问题描述】:
我们目前有一个 Python Apache Beam 管道正在工作并且能够在本地运行。我们现在正在让管道在 Google Cloud Dataflow 上运行并实现完全自动化,但发现 Dataflow/Apache Beam 的管道监控存在局限性。
目前,Cloud Dataflow 有两种监控管道状态的方法,一种是通过其 UI 界面,另一种是通过命令行中的 gcloud。这两种解决方案都不适用于我们可以考虑无损文件处理的全自动解决方案。
查看 Apache Beam 的 github,他们有一个文件 internal/apiclient.py,其中显示有一个用于获取作业状态的函数 get_job。
我们发现 get_job 使用的一个实例位于runners/dataflow_runner.py。
最终目标是使用此 API 获取我们自动触发运行的一个或多个作业的状态,以确保它们最终都通过管道成功处理。
谁能向我们解释在我们运行管道后如何使用此 API (p.run())?我们不明白response = runner.dataflow_client.get_job(job_id) 中的runner 来自哪里。
如果有人能够更深入地了解我们如何在设置/运行我们的管道时访问此 API 调用,那就太好了!
【问题讨论】:
标签: python pipeline google-cloud-dataflow apache-beam