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

如何在Airflow 2.3中调度串行执行的动态Task Group

解决方案:Airflow 2.3 动态串行Task Group实现

针对你的需求,Airflow 2.3已支持Task Group的动态映射,结合并行度控制可实现同一时间仅执行一个Task Group的效果,具体实现如下:

核心思路

  1. 通过上游Operator获取项目列表
  2. 定义Task Group工厂函数,内部维护两个Operator的顺序依赖
  3. 使用TaskGroup.expand()动态生成每个项目对应的Task Group,并设置parallelism=1强制串行执行

完整代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.task_group import TaskGroup
from datetime import datetime

# 模拟获取项目列表逻辑(实际可替换为从数据库/文件读取)
def fetch_project_list():
    return ["proj_001", "proj_002", "proj_003"]

# 定义Task Group工厂函数,每个项目对应一组顺序执行的任务
def build_project_task_group(project_id):
    with TaskGroup(group_id=f"task_group_{project_id}", tooltip=f"Pipeline for {project_id}") as task_group:
        # 第一个任务
        step_1 = PythonOperator(
            task_id="data_processing",
            python_callable=lambda proj: print(f"Processing {proj}: Step 1 completed"),
            op_args=[project_id]
        )
        # 第二个任务,依赖第一个任务完成
        step_2 = PythonOperator(
            task_id="result_upload",
            python_callable=lambda proj: print(f"Processing {proj}: Step 2 completed"),
            op_args=[project_id]
        )
        step_1 >> step_2
    return task_group

with DAG(
    dag_id="serial_dynamic_task_groups_dag",
    start_date=datetime(2023, 1, 1),
    schedule_interval=None,
    catchup=False,
    tags=["dynamic_task", "task_group"]
) as dag:
    # 步骤1:获取项目列表
    get_projects = PythonOperator(
        task_id="get_project_list",
        python_callable=fetch_project_list
    )

    # 步骤2:动态生成Task Group,parallelism=1确保同一时间仅执行一组
    dynamic_task_groups = build_project_task_group.expand(
        project_id=get_projects.output,
        parallelism=1  # 关键参数:限制并行执行的Task Group数量
    )

    # 设置依赖关系:先获取列表,再执行所有动态Task Group
    get_projects >> dynamic_task_groups

关键细节说明

  • 动态Task Group支持:Airflow 2.3及以上版本允许对TaskGroup对象使用.expand()方法,实现基于上游输出的动态生成
  • 串行控制:expand()方法的parallelism=1参数会限制同时运行的动态Task Group数量为1,完全满足"同一时间仅执行一个Task Group"的要求
  • 内部任务顺序:每个Task Group内部通过step_1 >> step_2明确了两个Operator的执行顺序,确保每组任务按流程完成

替代方案(非动态场景)

若坚持用遍历方式生成Task Group,可通过以下方式控制串行:

  1. 在DAG定义中设置max_active_tasks=1,强制整个DAG同一时间仅运行一个任务
  2. 为每个生成的Task Group单独设置concurrency=1,并手动维护Task Group之间的依赖链(前一个Task Group完成后再启动下一个)

但动态映射的方式更简洁,且天然适配"从上游获取项目列表"的场景,是更推荐的实现方式。

内容的提问来源于stack exchange,提问作者Pantonaut

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 18:05:27