【问题标题】:Publishing message to Pub/Sub after Dataflow template job done在 Dataflow 模板作业完成后将消息发布到 Pub/Sub
【发布时间】:2020-01-03 00:46:37
【问题描述】:

我想在我的 Dataflow 作业成功完成后向 Pub/Sub 主题发布消息。

我通过使用wait_until_finish() 来执行此操作,然后在我确认管道已完成后发布一条消息(当前状态为JOB_STATE_DONE)。当我在本地执行管道时(使用DirectRunnerDataflowRunner),我得到了想要的结果。但是,当我将此管道部署到模板然后从模板创建作业时,p.run() 之外的代码不会被执行。

我的管道到此结束

    result = p.run()
    result.wait_until_finish()
    publish_job_completion(result, all_pipeline_options)

publish_job_completion 确认工作已完成并发布消息。

我如何设计一个在完成后将发布到 Pub/Sub 的管道,即使是从模板运行时也是如此?

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    我不确定是否有通过模板管道本身执行此操作的好方法。但是,您可以将其包装到在 Cloud Composer 上运行的简单 Airflow DAG 中。 Airflow 有一个 DataflowTemplateOperator 可以启动数据流模板作业,还有一个 PubsubPublishOperator 可以将消息发布到 Pub/Sub。您可以将这两个操作串在一起,作为部署到 Cloud Composer 的简单 DAG。

    虽然this example 不使用模板化 Dataflow 作业,但它提供了一个很好的示例,说明如何使用 Cloud Composer 来编排 Dataflow 作业。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-01-01
      • 2021-06-23
      • 2021-08-01
      • 1970-01-01
      • 2022-08-22
      • 2022-08-17
      • 2019-07-23
      • 1970-01-01
      相关资源
      最近更新 更多