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

Airflow 2.3.x:如何用XCom配置TriggerDagRunOperator的max_active_tis_per_dag

解决方案:动态配置TriggerDagRunOperator的并行执行数

针对你需要从DAG运行配置中动态获取max_active_tis_per_dag参数的需求,以下是两种适用于Airflow 2.3.x的可行方案:

方案一:利用TaskFlow API自动传递参数(推荐)

TaskFlow API会自动处理任务间的依赖和数据传递,代码更简洁直观:

from airflow import DAG
from airflow.decorators import task
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from datetime import datetime

with DAG(
        'aaa_test_controller',
        schedule_interval=None,
        start_date=datetime(2021, 1, 1),
        catchup=False
) as dag:

    @task
    def get_num_max_parallel_runs(dag_run=None):
        # 从DAG运行配置中读取并行数,未配置时默认返回1
        return dag_run.conf.get("num_max_parallel_runs", 1)

    # 获取动态并行数,TaskFlow自动创建依赖关系
    parallel_run_count = get_num_max_parallel_runs()

    # 批量触发目标DAG,将动态并行数传入参数
    trigger_dag = TriggerDagRunOperator.partial(
        task_id="trigger_dependent_dag",
        trigger_dag_id="aaa_some_other_dag",
        wait_for_completion=True,
        max_active_tis_per_dag=parallel_run_count,  # 直接传入任务返回值
        poke_interval=5
    ).expand(conf=['{"some_key": "some_value_1"}', '{"some_key": "some_value_2"}'])

说明:

  • TaskFlow会自动在get_num_max_parallel_runs和trigger_dag之间建立执行依赖,确保先获取并行数再执行触发任务。
  • 任务返回值通过XCom机制自动传递,无需手动处理数据拉取逻辑。

方案二:使用模板变量手动拉取XCom值

如果遇到特殊场景导致方案一失效,可以用Jinja2模板语法直接拉取XCom数据:

from airflow import DAG
from airflow.decorators import task
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from datetime import datetime

with DAG(
        'aaa_test_controller',
        schedule_interval=None,
        start_date=datetime(2021, 1, 1),
        catchup=False
) as dag:

    @task
    def get_num_max_parallel_runs(dag_run=None):
        return dag_run.conf.get("num_max_parallel_runs", 1)

    # 定义获取并行数的任务
    get_parallel_task = get_num_max_parallel_runs()

    # 通过模板变量拉取XCom值,开启原生对象渲染保证类型正确
    trigger_dag = TriggerDagRunOperator.partial(
        task_id="trigger_dependent_dag",
        trigger_dag_id="aaa_some_other_dag",
        wait_for_completion=True,
        max_active_tis_per_dag="{{ ti.xcom_pull(task_ids='get_num_max_parallel_runs') }}",
        poke_interval=5,
        render_template_as_native_obj=True  # 确保解析后为整数类型,避免字符串错误
    ).expand(conf=['{"some_key": "some_value_1"}', '{"some_key": "some_value_2"}'])

    # 手动设置任务执行顺序
    get_parallel_task >> trigger_dag

说明:

  • {{ ti.xcom_pull(task_ids='get_num_max_parallel_runs') }}是Airflow模板语法,用于从指定任务拉取XCom数据。
  • render_template_as_native_obj=True避免模板解析后返回字符串类型,保证参数类型匹配。

使用注意:

  • 触发DAG时,可在Airflow UI的「Trigger DAG w/ config」中传入配置,例如{"num_max_parallel_runs": 2},未传入时将使用默认值1。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 05:40:30