Airflow:动态任务映射下实现分支内任务串行运行
动态任务分支串行+多分支并行实现方案
需求流程:
Task1(输出字典列表)→ 每个字典对应一条分支,分支内Task2与Task3串行执行,多分支可并行
基于Python标准库的实现
定义任务函数
def task1(): # 模拟输出字典列表 return [{"id": "a", "data": "content_a"}, {"id": "b", "data": "content_b"}] def task2(item): # 处理单个字典,模拟Task2逻辑 print(f"Executing Task2 for {item['id']}") return f"task2_result_{item['id']}" def task3(task2_result): # 接收Task2的结果,模拟Task3逻辑 print(f"Executing Task3 with result: {task2_result}") return f"task3_result_{task2_result.split('_')[-1]}"
实现分支串行+多分支并行
from concurrent.futures import ThreadPoolExecutor, as_completed def serial_task(item): # 单分支内Task2与Task3串行执行 t2_res = task2(item) t3_res = task3(t2_res) return t3_res def main(): # 执行Task1获取字典列表 items = task1() # 多分支并行处理 with ThreadPoolExecutor(max_workers=2) as executor: futures = [executor.submit(serial_task, item) for item in items] # 收集各分支结果 for future in as_completed(futures): print(f"Branch completed with result: {future.result()}") if __name__ == "__main__": main()
关键说明
serial_task封装单分支内的串行逻辑,确保Task2执行完成后才启动Task3ThreadPoolExecutor实现多分支并行,可通过max_workers参数控制并行数量- 若为CPU密集型任务,可替换为
ProcessPoolExecutor提升效率
基于工作流引擎(Airflow)的实现
针对调度类场景,可使用Airflow的动态任务映射能力:
from airflow.decorators import dag, task from datetime import datetime @dag(start_date=datetime(2024,1,1), schedule=None) def dynamic_branch_dag(): @task def task1(): return [{"id": "a", "data": "content_a"}, {"id": "b", "data": "content_b"}] @task def task2(item): print(f"Executing Task2 for {item['id']}") return f"task2_result_{item['id']}" @task def task3(task2_result): print(f"Executing Task3 with result: {task2_result}") return f"task3_result_{task2_result.split('_')[-1]}" # 动态生成分支,自动保证分支内串行 items = task1() t2_results = task2.expand(item=items) t3_results = task3.expand(task2_result=t2_results) dag = dynamic_branch_dag()
关键说明
expand方法根据Task1的输出动态生成对应数量的Task2实例- Airflow默认保证同一分支内Task2执行完成后才触发Task3,无需额外配置
- 多分支并行度可通过Airflow的全局配置或任务级参数调整
内容的提问来源于stack exchange,提问作者Shreya Singhal
相关产品推荐
相关产品推荐

