【问题标题】:How to launch a Dataflow job with Apache Airflow and not block other tasks?如何使用 Apache Airflow 启动 Dataflow 作业而不阻止其他任务?
【发布时间】:2020-01-25 14:59:03
【问题描述】:

问题

DataflowTemplateOperator 类型的气流任务需要很长时间才能完成。这意味着其他任务可以被它阻止(对吗?)。

当我们运行更多这样的任务时,这意味着我们需要更大的 Cloud Composer 集群(在我们的例子中)来执行本质上是阻塞而它们不应该阻塞的任务(它们应该是 async 操作)。

选项

  • 选项 1:只需启动作业,气流作业就成功了
  • 选项 2:按照 here 的说明编写包装器,并按照 here 的说明使用重新调度模式

选项 1 似乎不可行,因为 DataflowTemplateOperator 只有一个选项来指定完成检查之间的等待时间,称为 poll_sleep (source)。

对于DataflowCreateJavaJobOperator,有一个选项check_if_running 可以等待完成以前的同名作业 (see this code)

似乎在启动一个作业后,wait_for_finish 被执行(参见this line),这归结为一个“不完整”的作业(参见this line)。

对于选项 2,我需要选项 1。

问题

  1. 我认为 Dataflow 任务会阻止 Cloud Composer/Airflow 中的其他任务是否正确?
  2. 有没有办法使用内置运算符无需“等待完成”来安排作业? (我可能忽略了一些东西)
  3. 有没有一种简单的方法可以自己编写?我正在考虑只执行一个 bash 启动脚本,然后执行一个任务来查看作业是否正确完成,但处于重新安排模式。
  4. 在运行数据流作业时是否有另一种避免阻塞其他任务的方法?基本上这是一个异步操作,不应该占用资源。

【问题讨论】:

    标签: airflow google-cloud-dataflow apache-beam google-cloud-composer


    【解决方案1】:

    答案

    1. 我认为 Dataflow 任务会阻止 Cloud Composer/Airflow 中的其他任务是否正确?
      答:部分是的。 Airflow 在配置中具有并行选项,它定义了一次应在整个系统中执行的任务数量。让任务阻塞此插槽可能会减慢系统中的执行速度,但随着您增加任务和 DAG 的数量,这个问题必然会发生。您可以根据需要在配置中增加此项

    1. 有没有一种方法可以使用内置运算符来安排作业而无需“等待完成”? (我可能忽略了一些东西)
      答:是的。您可以使用 PythonOperator 并在 python_callable 中使用数据流挂钩以异步模式启动作业(启动且不要等待)。

    1. 有没有一种简单的方法可以自己编写?我正在考虑只执行一个 bash 启动脚本,然后执行一个任务,该任务查看作业是否正确完成,但处于重新安排模式。 答:当您说重新安排时,我假设您将重试查找作业以检查作业是否正确完成的任务。如果我是对的,您可以将任务设置为重试模式以及您希望重试发生的延迟时间。

    1. 是否有其他方法可以避免在运行数据流作业时阻塞其他任务?基本上这是一个异步操作,不应该占用资源。
      A:我想我在第二个问题中回答了这个问题。

    【讨论】:

    • 谢谢!知道如何以异步模式启动任务吗?
    • 同步或异步完全取决于您的依赖流...当您添加对数据流模板运算符的依赖时,下游任务将始终依赖于数据流运算符...如果您不想要这个,您必须设计一个依赖项,以确保作业以异步模式启动并且 Airflow 任务不会被阻止......
    • 如果等待数据流模板运算符的时间过长给您带来了问题,您可以创建自己的数据流模板运算符,当数据流作业未处于其最终状态(JOB_STATE_DONE 或JOB_STATE_FAILED) 并标记一定延迟后重试...这也可以确保任务不被阻塞并且您可以在同步模式下进行任务...
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-03-07
    • 2017-09-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-13
    相关资源
    最近更新 更多