Airflow如何实现解析时未知数量的动态顺序任务组?
Airflow 动态顺序批次任务组实现方案
实现思路
基于Airflow原生功能,用TaskFlow API实现运行时分批,动态生成Task Group后通过链式依赖让批次按顺序执行,批次内任务并行处理,完全无需硬编码批次数量。
完整代码示例
from airflow import DAG from airflow.decorators import task, task_group from airflow.models import Variable from airflow.utils.helpers import chain from datetime import datetime # 替换为你的业务处理逻辑 def process_item_logic(item): print(f"处理项目: {item}") return item @task def fetch_project_list(): # 从Airflow Variable读取项目列表 return Variable.get("project_list", deserialize_json=True) @task def split_to_batches(project_list, batch_size=3): # 运行时分批,可自定义批次大小 batches = [] for idx in range(0, len(project_list), batch_size): batches.append(project_list[idx:idx+batch_size]) return batches @task_group def batch_processing_group(batch): # 单个批次的处理组,内部任务并行 @task def process_item(item): return process_item_logic(item) # 动态映射当前批次的所有项目,生成并行任务 process_item.expand(item=batch) with DAG( dag_id="dynamic_sequential_batch_dag", start_date=datetime(2024, 1, 1), schedule=None, catchup=False ) as dag: # 第一步:获取项目列表并分批 project_list = fetch_project_list() task_batches = split_to_batches(project_list) # 动态生成所有批次的任务组,并按顺序链式连接 batch_task_groups = [] for batch in task_batches: group = batch_processing_group(batch=batch) batch_task_groups.append(group) # 确保批次按顺序执行:上一批全完成才启动下一批 chain(*batch_task_groups)
核心细节说明
- 运行时动态性:批次的数量和内容完全由
split_to_batches任务在DAG运行时生成,不需要在DAG解析阶段提前确定,规避了硬编码的问题。 - 顺序执行保障:
chain(*batch_task_groups)会把所有生成的Task Group按顺序串联,确保前一个批次的所有任务执行完毕后,才会启动下一个批次的任务。 - 批次内并行:每个Task Group内部用
.expand()实现批次内项目的并行处理,兼顾了执行效率和顺序要求。 - 版本要求:该方案适用于Airflow 2.3及以上版本(2.3+支持TaskFlow的动态任务映射,2.5+对动态Task Group的支持更稳定)。
低版本替代方案
如果使用Airflow 2.3以下版本,可采用TriggerDagRunOperator触发子DAG的方式,每个批次对应一个子DAG实例,通过父DAG的任务依赖控制批次顺序,但这种方式维护成本较高,官方更推荐升级到2.3+版本使用上述原生方案。
内容的提问来源于stack exchange,提问作者Aashish Vinu
相关产品推荐
相关产品推荐

