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
关键逻辑说明
- 并行启动:
task_a1和task_b1没有设置上游任务,会在DAG启动时同时执行。 - 串行分支:通过
>>操作符依次串联A1→A2、B1→B2→B3的串行关系。 - 多上游等待:把
task_a2和task_b3放在列表里作为task_c1的上游,Airflow会自动等待这两个任务都成功完成后,才触发task_c1执行。 - 后续串行:
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
相关产品推荐
相关产品推荐

