【问题标题】:Rerun Airflow Dag From Middle Task and continue to run it till end of all downstream tasks.(Resume Airflow DAG from any task)从中间任务重新运行 Airflow Dag 并继续运行它直到所有下游任务结束。(从任何任务恢复 Airflow DAG)
【发布时间】:2018-10-10 16:39:57
【问题描述】:

嗨,我是 Apache Airflow 的新手,可以说我有很多依赖关系

任务 A >> 任务 B >> 任务 C >> 任务 D >> 任务 E

  1. 是否可以从中间任务运行 Airflow DAG,比如说任务 C?

  2. 在分支的情况下是否可以只运行特定的分支 中间的算子?

  3. 是否可以从上次失败的任务中恢复 Airflow DAG?

  4. 如果不可能,如何管理大型 DAG 并避免重新运行 多余的任务?

如果可能,请向我提供有关如何实施的建议。

【问题讨论】:

    标签: airflow


    【解决方案1】:
    1. 您不能手动执行此操作。如果您设置 BranchPythonOperator,您可以根据 BranchPythonOperator 中设置的条件跳过任务,直到您希望开始的任务

    2. 同1。

    3. 是的。您可以清除上游任务直到根节点或下游直到节点的所有叶子。

    你可以这样做:

    Task A >> Task B >> Task C >> Task D
    Task C >> Task E
    

    其中C 是分支运算符。 例如:

        from datetime import date
        def branch_func():
            if date.today().weekday() == 0:
                return 'task id of D'
            else:
                return 'task id of E'
    
    
        Task_C = BranchPythonOperator(
            task_id='branch_operation',
            python_callable=branch_func,
            dag=dag)
    

    这将是周一的任务序列:

    Task A >> Task B >> Task C >> Task D
    

    这将是本周剩余时间的任务序列:

    Task A >> Task B >> Task C >> Task E
    

    【讨论】:

    • 是的分支运算符存在,但我的问题是从中间开始 dag 。在分支情况下,也可以从分支启动它,即任务 C
    • 您可以从超过 1 个点开始 DAG,但您必须设置条件以跳过您希望不运行的分支。
    • 是的,我尝试使用分支并跳过任务,但是当我只触发分支任务时,它不会从分支继续到结束。在上面的示例中,正如您提到的,如果我点击命令,例如“airflow run dag_id task_c date”,那么在我的 UI 中,我可以看到 task_c 正在执行 task_d,但是如果我在 task_d 之后还有更多任务,可以说 task_f 它不工作。你有一些例子或博客可以分享一下吗?
    • 为什么不分离多个 DAG?
    • 是的,这是我们的选择之一。所以我认为没有任何分支直接做它是不可能的。
    猜你喜欢
    • 2021-06-19
    • 1970-01-01
    • 2019-08-12
    • 1970-01-01
    • 1970-01-01
    • 2020-06-21
    • 1970-01-01
    • 2022-11-11
    • 1970-01-01
    相关资源
    最近更新 更多