如何设置Airflow 2.2.2的rollups DAG每次仅运行一个任务?
问题描述
使用Python3与Airflow 2.2.2搭建数据管道,主DAG「rollups」包含大量针对Postgres实例的SQL查询任务。由于任务间无依赖关系(大量兄弟任务),Airflow默认会同时启动所有满足条件的任务,导致5-10个任务并行运行,使小型Postgres实例过载锁死。
当前通过将所有任务设置为顺序依赖(B << A、C << B等)实现单任务运行,但破坏了任务间的逻辑依赖结构,形成过深的任务树。希望在保留合理依赖结构的同时,强制该DAG每次仅运行一个任务,但设置concurrency和max_active_tasks_per_dag等参数后未解决问题,寻求可行方案。
解决方案
方案1:使用任务池(Task Pool)
Airflow的任务池机制可限制特定任务组的并发数。为「rollups」DAG的所有任务分配一个容量为1的专属任务池,即可确保同一时间仅一个任务运行,同时保留任务原有的依赖关系。
步骤:
- 在Airflow UI的Admin > Pools页面创建新池,例如命名为
postgres_rollups_pool,设置Slots为1。 - 在DAG中定义任务时,为每个任务指定
pool参数:
from airflow import DAG from airflow.operators.postgres_operator import PostgresOperator from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 1, 1), } with DAG('rollups', default_args=default_args, schedule_interval='@daily') as dag: task1 = PostgresOperator( task_id='task1_rollup', sql='SELECT * FROM rollup_table1;', postgres_conn_id='postgres_default', pool='postgres_rollups_pool' # 指定任务池 ) task2 = PostgresOperator( task_id='task2_rollup', sql='SELECT * FROM rollup_table2;', postgres_conn_id='postgres_default', pool='postgres_rollups_pool' ) task3 = PostgresOperator( task_id='task3_rollup', sql='SELECT * FROM rollup_table3;', postgres_conn_id='postgres_default', pool='postgres_rollups_pool' ) # 保留原有的依赖关系(如果有) task3 << task1
方案2:正确配置DAG级并发参数
Airflow 2.2.2中,控制单个DAG实例并发任务数的核心参数是DAG级的concurrency,而非全局的max_active_tasks_per_dag。需确保在DAG定义中明确设置该参数,且全局配置未覆盖DAG级设置。
代码示例:
from airflow import DAG from airflow.operators.postgres_operator import PostgresOperator from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 1, 1), } # 设置DAG的concurrency为1,限制该DAG同一时间仅运行1个任务 with DAG('rollups', default_args=default_args, schedule_interval='@daily', concurrency=1, # 关键参数 max_active_runs=1) as dag: # 可选:限制该DAG同时仅1个活跃运行实例 task1 = PostgresOperator( task_id='task1_rollup', sql='SELECT * FROM rollup_table1;', postgres_conn_id='postgres_default' ) task2 = PostgresOperator( task_id='task2_rollup', sql='SELECT * FROM rollup_table2;', postgres_conn_id='postgres_default' ) task3 = PostgresOperator( task_id='task3_rollup', sql='SELECT * FROM rollup_table3;', postgres_conn_id='postgres_default' ) # 保留原依赖结构 task3 << task1
注意事项:
- 若全局配置中的
core.dag_concurrency值小于DAG级concurrency,则全局配置会生效,需确保全局配置未限制过严。 max_active_runs控制同一DAG的活跃运行实例数,若需确保单实例单任务,可同时设置为1。
方案3:使用TaskGroup限制并发
如果仅需限制某一组无依赖任务的并发数,可将这些任务放入TaskGroup,并设置TaskGroup的concurrency参数为1,既保留任务间的逻辑分组,又控制并发。
代码示例:
from airflow import DAG from airflow.operators.postgres_operator import PostgresOperator from airflow.utils.task_group import TaskGroup from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 1, 1), } with DAG('rollups', default_args=default_args, schedule_interval='@daily') as dag: # 前置任务(如果有) pre_task = PostgresOperator( task_id='pre_rollup_check', sql='SELECT COUNT(*) FROM source_table;', postgres_conn_id='postgres_default' ) # 创建TaskGroup,设置concurrency=1 with TaskGroup('rollup_tasks', concurrency=1) as rollup_group: task1 = PostgresOperator( task_id='task1_rollup', sql='SELECT * FROM rollup_table1;', postgres_conn_id='postgres_default' ) task2 = PostgresOperator( task_id='task2_rollup', sql='SELECT * FROM rollup_table2;', postgres_conn_id='postgres_default' ) task3 = PostgresOperator( task_id='task3_rollup', sql='SELECT * FROM rollup_table3;', postgres_conn_id='postgres_default' ) # 保留依赖逻辑:前置任务完成后,执行rollup组内的任务(组内并发为1) pre_task >> rollup_group
内容的提问来源于stack exchange,提问作者James Webb
相关产品推荐
相关产品推荐

