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

如何通过Apache Airflow向Azure Data Factory管道传递运行参数

问题描述

我尝试在Apache Airflow中通过手动触发并携带配置的方式,将运行参数传递给Azure Data Factory管道。已配置好可通过运行时配置更新参数,但在将这些参数传入pipeline_parameters指定位置时遇到问题。我原本以为参数应作为字典存入args中,因此尝试直接从args['params']取值分配到pipeline_parameters,以下是简化后的代码:

import configparser
import sys
import pendulum

from os import path
from datetime import timedelta, datetime
from airflow import DAG
from airflow.models import Variable
from airflow.operators.dummy import DummyOperator
from airflow.utils.dates import days_ago
from airflow import DAG, Dataset

from custom_interfaces.alerts import send_alert
from custom_interfaces.models.message import MessageSettings
from custom_interfaces import AzureDataFactoryOperator

args = {
    "owner": "user-team-hr-001",
    "start_date": datetime(2021, 1, 1, tzinfo=local_tz),
    "on_failure_callback": on_failure_callback,
    "retries": 3,
    'retry_delay': timedelta(seconds=20),
    'params': {
        "p_date_overwrite":"",
        "p_foldername":"",
        "p_foler_prefix":""
    }
}

with DAG(
    dag_id=DAG_ID,
    description="Run with Pipeline Parameters.",
    catchup=False,
    default_args = args,
    dagrun_timeout=timedelta(minutes=20),
    is_paused_upon_creation=True
  
) as dag:

    run_pipeline_workday_bpt_historic_load = AzureDataFactoryOperator(
        task_id="test_load",
        trigger_rule="all_done",
        adf_sp_connection_id=config['adf_sp_connection_id'],
        subscription_id=config['subscriptions']['id'],
        resource_group_name=config['adf_resource_group'],
        factory_name=config['adf_name'],
        pipeline_name=config['adf_pipeline_name'],
        pipeline_parameters={
            "p_date_overwrite":args['params']['p_date_overwrite'],
            "p_foldername":args['params']['p_foldername'],
            "p_foler_prefix":args['params']['p_foler_prefix']
        },
        polling_period_seconds=10,
        outlets=[dataset_name]
    )
解决方案

你遇到的问题核心在于:args['params']里存储的是参数的默认值,手动触发DAG时传入的运行时参数不会更新这个字典,而是存储在DAG运行的上下文(Context)中。因此需要通过Airflow的Jinja模板语法来动态获取运行时参数值。

修改AzureDataFactoryOperator中的pipeline_parameters配置,使用Jinja模板变量引用参数:

run_pipeline_workday_bpt_historic_load = AzureDataFactoryOperator(
    task_id="test_load",
    trigger_rule="all_done",
    adf_sp_connection_id=config['adf_sp_connection_id'],
    subscription_id=config['subscriptions']['id'],
    resource_group_name=config['adf_resource_group'],
    factory_name=config['adf_name'],
    pipeline_name=config['adf_pipeline_name'],
    # 使用Jinja模板引用运行时参数
    pipeline_parameters={
        "p_date_overwrite": "{{ params.p_date_overwrite }}",
        "p_foldername": "{{ params.p_foldername }}",
        "p_foler_prefix": "{{ params.p_foler_prefix }}"
    },
    polling_period_seconds=10,
    outlets=[dataset_name]
)

补充说明

  • 确保你的AzureDataFactoryOperator支持模板化pipeline_parameters参数(Airflow官方或主流自定义实现均支持)。
  • 手动触发DAG时,在"配置"栏填写对应的参数键值对,比如:
    {
        "p_date_overwrite": "2024-05-20",
        "p_foldername": "hr_data",
        "p_foler_prefix": "workday_"
    }
    
  • 如果需要处理空值情况,可以在模板中添加默认值,例如:{{ params.p_date_overwrite or 'default_date' }}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 09:09:55