如何通过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
相关产品推荐
相关产品推荐

