如何在Airflow触发DAG时基于传入参数解析Operator的连接ID?
动态传递触发参数切换FTP连接实现FTP到Azure Blob传输
要实现触发DAG时指定FTP连接,核心是利用Airflow的Jinja模板机制引用运行时传入的配置参数,直接将参数值注入Operator的对应字段。以下是具体实现方案:
核心思路
Airflow允许在Operator的支持模板化的字段中使用dag_run.conf变量获取触发DAG时传入的参数(比如通过UI触发时填写的JSON配置)。大部分连接ID类字段(如ftp_conn_id)默认支持模板渲染,无需额外配置。
修改后的DAG代码
from airflow import DAG from airflow.providers.ftp.operators.ftp import FTPFileTransmitOperator, FTPOperation from airflow.providers.microsoft.azure.operators.wasb import LocalFilesystemToWasbOperator from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1) } with DAG( "example_ftp_to_blob", default_args=default_args, schedule=None, catchup=False, ) as dag: ftp_get = FTPFileTransmitOperator( task_id="ftp_get", # 从触发参数中获取FTP连接ID,设置默认值避免无参数时出错 ftp_conn_id="{{ dag_run.conf.get('ftp_conn_id', 'my-ftp-con') }}", # 可选:文件路径也可以通过参数动态指定 local_filepath="{{ dag_run.conf.get('local_filepath', '/tmp/my-file') }}", remote_filepath="{{ dag_run.conf.get('remote_filepath', '/my-file') }}", operation=FTPOperation.GET, create_intermediate_dirs=True, ) blob_put = LocalFilesystemToWasbOperator( task_id="blob_put", wasb_conn_id="my_blob_conn", # 同步FTP任务的本地文件路径 file_path="{{ dag_run.conf.get('local_filepath', '/tmp/my-file') }}", container_name="incoming", blob_name="{{ dag_run.conf.get('blob_name', 'my-file') }}", ) ftp_get >> blob_put
使用方法
触发DAG时,在Airflow UI的配置JSON框中传入指定的参数,示例:
{ "ftp_conn_id": "my-ftp-con-2", "remote_filepath": "/path/to/another-file", "blob_name": "another-file" }
注意事项
- 确保Airflow已配置好对应的FTP连接(如
my-ftp-con-1、my-ftp-con-2),连接信息在Airflow UI的Admin > Connections中维护。 - 如果需要对参数做额外校验或处理,可以在DAG中添加
PythonOperator前置任务,检查dag_run.conf的合法性,避免无效参数导致任务失败。
内容的提问来源于stack exchange,提问作者Christopher Bennage
相关产品推荐
相关产品推荐

