You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow任务并行与依赖配置咨询:如何实现指定调度结构?

Airflow DAG 混合并行/串行任务依赖配置

针对你描述的任务流,直接用Airflow的任务依赖操作符就能轻松实现,核心是利用多任务作为上游的特性来实现等待所有分支完成的逻辑。

完整代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

# 定义基础任务函数(可替换为你实际的任务逻辑)
def task_func(task_name):
    print(f"Executing task: {task_name}")

# 初始化DAG
with DAG(
    dag_id="mixed_parallel_serial_dag",
    start_date=datetime(2024, 1, 1),
    schedule_interval="@daily",
    catchup=False
) as dag:
    # 定义所有任务
    task_a1 = PythonOperator(
        task_id="A1",
        python_callable=task_func,
        op_kwargs={"task_name": "A1"}
    )

    task_b1 = PythonOperator(
        task_id="B1",
        python_callable=task_func,
        op_kwargs={"task_name": "B1"}
    )

    task_a2 = PythonOperator(
        task_id="A2",
        python_callable=task_func,
        op_kwargs={"task_name": "A2"}
    )

    task_b2 = PythonOperator(
        task_id="B2",
        python_callable=task_func,
        op_kwargs={"task_name": "B2"}
    )

    task_b3 = PythonOperator(
        task_id="B3",
        python_callable=task_func,
        op_kwargs={"task_name": "B3"}
    )

    task_c1 = PythonOperator(
        task_id="C1",
        python_callable=task_func,
        op_kwargs={"task_name": "C1"}
    )

    task_c2 = PythonOperator(
        task_id="C2",
        python_callable=task_func,
        op_kwargs={"task_name": "C2"}
    )

    # 配置任务依赖
    # A分支:A1 → A2
    task_a1 >> task_a2
    # B分支:B1 → B2 → B3
    task_b1 >> task_b2 >> task_b3
    # 等待A2和B3都完成后启动C1,再到C2
    [task_a2, task_b3] >> task_c1 >> task_c2

关键逻辑说明

  1. 并行启动:task_a1和task_b1没有设置上游任务,会在DAG启动时同时执行。
  2. 串行分支:通过>>操作符依次串联A1→A2、B1→B2→B3的串行关系。
  3. 多上游等待:把task_a2和task_b3放在列表里作为task_c1的上游,Airflow会自动等待这两个任务都成功完成后,才触发task_c1执行。
  4. 后续串行:task_c1完成后直接触发task_c2,保持串行逻辑。

替代写法(如果习惯用set_upstream)

如果你更倾向于用set_upstream方法,也可以这样配置依赖:

task_a2.set_upstream(task_a1)
task_b2.set_upstream(task_b1)
task_b3.set_upstream(task_b2)
task_c1.set_upstream([task_a2, task_b3])
task_c2.set_upstream(task_c1)

效果和用>>完全一致,选你觉得更直观的写法就行。

内容的提问来源于stack exchange,提问作者unnest_me

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 06:01:04