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

Airflow主DAG如何通过Postman传入参数更新子DAG执行次数

问题描述

主DAG需要根据Postman传入的JSON参数决定子DAG的执行次数,当前尝试通过全局变量更新执行次数时,遇到错误:airflow.exceptions.AirflowException: conf parameter should be JSON Serializable。

环境信息

  • Airflow版本:2.4.2
  • Postman触发请求Payload:
{
  "dag_run_id": "example_triger_several_time_a_sub_dag_001",
  "conf": {
    "steps": 20
  }
}

原代码问题分析

  1. DAG文件在调度器解析阶段就会执行for i in range(int(local_conf["steps"]))循环,此时还未触发DAG Run,无法获取到Postman传入的conf参数,循环次数只能用硬编码的13。
  2. update_globals函数中把整个context赋值给全局变量,context包含大量不可序列化的Airflow内部对象(如TaskInstance、Connection等),后续传递给conf时会触发序列化错误。
解决方案

利用Airflow 2.4+支持的**动态任务映射(Dynamic Task Mapping)**实现根据传入参数动态生成子DAG触发任务,步骤如下:

  1. 移除全局变量,直接从DAG Run的conf中读取steps参数。
  2. 使用TriggerDagRunOperator的动态映射功能,根据steps值生成对应数量的任务。
  3. 确保传递给子DAG的conf是纯JSON可序列化的字典。
修改后的代码
# main_dag.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from airflow.utils.dates import days_ago

''' Set default args '''
default_args = {
    'owner': 'rtm-pa-consumerization',
    'start_date': datetime(2021, 6, 25)
}

with DAG(
    dag_id='Example_trigger_sub_dag',
    default_args=default_args,
    schedule_interval=None,
    catchup=False
) as dag:
    
    ''' Declare check-points '''
    start = DummyOperator(task_id='Start')
    end = DummyOperator(task_id='End')
    
    ''' 动态生成子DAG触发任务 '''
    trigger_sub_dag = TriggerDagRunOperator.partial(
        task_id='Data_sub_dag_Q',
        trigger_dag_id='Sub_dag',
        reset_dag_run=True,
        wait_for_completion=True
    ).expand(
        # 从DAG Run的conf中获取steps,生成对应数量的任务
        conf=[
            {"process_id": i+1, "steps": "{{ dag_run.conf.get('steps', 13) }}"} 
            for i in range(int("{{ dag_run.conf.get('steps', 13) }}"))
        ]
    )

    ''' 任务依赖 '''
    start >> trigger_sub_dag >> end
关键说明
  • 动态映射:使用partial()定义基础算子配置,expand()根据dag_run.conf中的steps参数动态生成任务实例,每个任务的process_id自动递增。
  • 参数获取:通过{{ dag_run.conf.get('steps', 13) }}模板语法获取传入的参数,默认值设为13保证未传参时也能正常运行。
  • 序列化安全:传递给conf的是纯字典对象,没有包含Airflow内部不可序列化的对象,避免了之前的报错。
替代方案(如果无法使用动态映射)

如果因版本限制无法使用动态映射,可通过PythonOperator生成任务列表并使用XCom传递参数,再用BranchPythonOperator触发对应任务,但实现复杂度更高,推荐优先使用动态映射。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 22:50:25