【问题标题】:Python Apache Beam Pipeline Status API CallPython Apache Beam 管道状态 API 调用
【发布时间】: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


    【解决方案1】:

    我最终只是摆弄代码并找到了如何获取工作详细信息。我们的下一步是看看是否有办法获取所有工作的列表。

    # start the pipeline process
    pipeline                 = p.run()
    # get the job_id for the current pipeline and store it somewhere
    job_id                   = pipeline.job_id()
    # setup a job_version variable (either batch or streaming)
    job_version              = dataflow_runner.DataflowPipelineRunner.BATCH_ENVIRONMENT_MAJOR_VERSION
    # setup "runner" which is just a dictionary, I call it local
    local                    = {}
    # create a dataflow_client
    local['dataflow_client'] = apiclient.DataflowApplicationClient(pipeline_options, job_version)
    # get the job details from the dataflow_client
    print local['dataflow_client'].get_job(job_id)
    

    【讨论】:

    • 嘿@T.Okahara,你有什么改变想办法用数据流模板做到这一点吗?
    • 对不起@jmoore255 除了上面的代码之外,我们没有更多地让我们的管道在 Cloud Dataflow 中运行。我们实际上构建了自己的本地运行机器来运行我们的流程,因为我们发现在 Dataflow 上运行的其他问题,例如不允许我们从 App Engine 触发以及启动/清理时间缓慢。现在可能有所不同,但我们仍然在本地运行我们的管道(为 ML 进行数据修改)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-06-10
    • 1970-01-01
    • 1970-01-01
    • 2022-12-24
    • 1970-01-01
    • 2019-06-22
    • 2019-04-23
    相关资源
    最近更新 更多