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

Airflow中通过DAG运行配置实现动态任务序列编排

Airflow中通过DAG运行配置实现动态任务序列编排

我完全理解你的困境——硬编码的元素列表虽然能正常跑起来,但每次要调整任务序列都得改代码重新部署,实在太不灵活了!别着急,我来给你讲讲怎么用Airflow的dag_run.conf实现动态任务编排,完美解决你的问题。

核心思路

Airflow的DAG是静态解析的(调度器启动时就会解析DAG定义),这时候还没有具体的dag_run实例,所以没法直接在DAG定义阶段读取dag_run.conf。我们需要利用Airflow 2.x的**Dynamic Task Mapping(动态任务映射)**特性,在运行时根据dag_run.conf的参数动态生成任务和依赖关系。

具体实现步骤

1. 触发DAG时传入动态配置

手动触发DAG时,在「Configuration」栏传入JSON格式的elements参数,示例配置如下:

{
  "elements": [
    [["a", "b"], ["c", "d"], ["e", "f"], ["g", "h"], ["i", "j"]],
    [["k", "l"], ["m", "n"], ["o", "p"], ["q", "r"]],
    [["s", "t"], ["u", "v"], ["w", "x"]],
    [["y", "z"]]
  ]
}

2. 完整代码实现

from airflow import DAG
from airflow.decorators import task, task_group
from airflow.operators.python import get_current_context
from datetime import datetime
from airflow.utils.helpers import chain

@task
def process_element(element: str) -> str:
    """处理单个元素的基础任务"""
    print(f"Processing element: {element}")
    return f"Processed_{element}"

@task_group
def process_single_group(group_elements: list):
    """处理一组元素,组内任务并行执行"""
    # 生成组内所有处理任务
    group_tasks = [process_element.override(task_id=f"process_{elem}")(elem) for elem in group_elements]
    return group_tasks

@task
def fetch_elements() -> list:
    """从dag_run.conf中获取需要处理的elements列表"""
    context = get_current_context()
    dag_run_conf = context.get("dag_run", {}).conf or {}
    elements = dag_run_conf.get("elements", [])
    if not elements:
        raise ValueError("未在dag_run.conf中找到elements参数,请检查触发配置!")
    return elements

with DAG(
    dag_id="dynamic_task_sequence_from_conf",
    description="通过dag_run.conf动态生成任务序列与依赖的Airflow DAG",
    start_date=datetime(2025, 1, 28),
    schedule_interval=None,
    catchup=False,
    tags=["dynamic", "dag_run_conf"]
) as dag:
    # 第一步:获取运行时的elements列表
    elements_list = fetch_elements()
    
    # 第二步:定义处理单个任务序列的任务组(组间串行)
    @task_group
    def process_sequence(sequence: list):
        """处理一个完整的任务序列,组间串行执行"""
        prev_group = None
        for idx, group in enumerate(sequence):
            current_group = process_single_group.override(task_id=f"group_{idx}")(group)
            if prev_group:
                # 设置前一个组完成后再执行当前组
                prev_group >> current_group
            prev_group = current_group
        return prev_group
    
    # 动态生成所有序列的任务组
    sequence_groups = process_sequence.expand(sequence=elements_list)
    
    # 可选:让多个序列之间也串行执行(如果不需要可以删除此行)
    chain(*sequence_groups)

代码解释

  • fetch_elements任务:在DAG运行时通过Airflow上下文获取dag_run.conf中的elements列表,确保我们拿到动态配置的任务序列。
  • process_single_group任务组:负责处理每个子任务组,组内的process_element任务会并行执行。
  • process_sequence任务组:负责处理一个完整的任务序列(比如你原来elements中的第一个子列表),内部会自动设置组与组之间的串行依赖。
  • expand方法:根据elements_list动态生成所有序列的任务组,实现完全的动态任务编排。
  • chain方法:可选操作,用来让多个任务序列之间也串行执行,如果你的序列之间不需要依赖,可以删除这一行。

注意事项

  • 确保你的Airflow版本在2.3+,因为动态任务映射(expand)是在这个版本之后稳定支持的。
  • 触发DAG时必须传入elements参数,否则fetch_elements会抛出错误,你也可以给elements设置一个默认空列表来避免报错。
  • 任务ID会自动生成唯一值,避免重复,因为我们用了override(task_id=...)来设置个性化的任务ID。

备注:内容来源于stack exchange,提问作者james gem

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 15:38:09