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

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执行完成后才启动Task3
  • ThreadPoolExecutor实现多分支并行,可通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 13:00:46