【问题标题】:Airflow DAG Loop - How to make each iteration sequential instead of parallel气流 DAG 循环 - 如何使每个迭代顺序而不是并行
【发布时间】:2019-10-11 01:03:36
【问题描述】:

我有一个像这样的 Apache Airflow DAG:

DAG_NAME='my_dag'
sections = ["0", "1", "2", "3"]

with DAG(DAG_NAME, default_args=default_args, schedule_interval=None) as dag:

        for s in sections:
            a = DummyOperator(task_id=f"section_{s}_start")
            b = SubDagOperator(task_id=f"init_{s}_subdag",subdag=init_section(DAG_NAME,f"init_{s}_subdag", default_args))
            c = SubDagOperator(task_id=f"process_{s}_subdag", subdag=process_section(DAG_NAME,f"process_{s}_subdag", default_args))
            d = SubDagOperator(task_id=f"update_{s}_subdag", subdag=update_section(DAG_NAME,f"update_{s}_subdag", default_args))
            e = DummyOperator(task_id=f"section_{s}_end")
            a>>b>>c>>d>>e

此代码将我的任务呈现为 so

我怎样才能使任务的顺序是:

section_0_start>>init_0_subdag>>process_0_subdag>>update_0_subdag>>section_0_end section_0_end>>section_1_start section_1_start>>init_1_subdag>>process_1_subdag>>update_1_subdag>>section_1_end

.....

依此类推,从第 0 节开始,到第 3 节任务结束

谢谢

【问题讨论】:

  • 你确定需要 subDAG 吗?
  • @MeghdeepRay 每个 subdag 中有更多需要并行运行的任务。例如,我在每个处理 subdag 中读取 5 个需要并行运行的文件。此外,我希望我的团队为他们的代码重用我的 subdags,所以我试图制作一个通用模板。您还有其他更好的方法来实现这一点吗?

标签: airflow directed-acyclic-graphs


【解决方案1】:

像这样修改for循环:

    previous_e = None
    for s in sections:
        a = ...
        ...
        e = ...
        if previous_e:
            previous_e >> a
        a>>b>>c>>d>>e
        previous_e = e

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-02-24
    • 2022-07-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多