如何通过POST请求携带参数触发Airflow DAG?
触发MWAA DAG并传递参数的实现方案
咱们一步步来实现这个需求:在AWS MWAA上触发一个DAG,传递参数并让任务打印出这个参数值。
1. 编写接收参数的DAG
首先需要创建一个能接收触发参数的DAG,里面包含一个打印参数的任务。这里用PythonOperator,通过上下文获取触发时传入的参数:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime, timedelta def print_param(**context): # 从DAG运行上下文的配置中提取参数x,未传入时用默认值兜底 x_value = context['dag_run'].conf.get('x', '未传入参数') print(f"传入的参数x的值是: {x_value}") # DAG的默认参数配置 default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2023, 1, 1), 'email_on_failure': False, 'email_on_retry': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } with DAG( 'test_param_dag', default_args=default_args, description='测试接收触发参数的DAG', schedule_interval=None, # 设置为None,仅支持手动触发 catchup=False, ) as dag: print_task = PythonOperator( task_id='print_parameter', python_callable=print_param, provide_context=True, # 必须开启这个参数,才能让函数获取Airflow上下文 ) print_task
关键说明
provide_context=True:开启后,Python函数才能拿到Airflow的上下文对象,进而读取DAG运行时的配置参数。context['dag_run'].conf:这个对象存储了触发DAG时传入的所有参数,用get()方法可以安全提取指定参数,还能设置默认值避免报错。
2. 编写触发DAG的Python代码
接下来用Python代码调用MWAA的CLI接口触发DAG,并传递参数。这里要注意:传递给DAG的参数需要放在CLI命令的--conf参数中,用JSON字符串格式传递,而不是HTTP请求的查询参数。
import base64 import boto3 import requests MWAA_ENVIRONMENT_NAME = 'production-mwaa' dag_name = 'test_param_dag' mwaa_cli_command = 'dags trigger' # 定义要传递给DAG的参数,转为JSON字符串格式 dag_conf = '{"x": 11}' # 创建MWAA客户端并获取CLI访问token client = boto3.client('mwaa') mwaa_cli_token = client.create_cli_token(Name=MWAA_ENVIRONMENT_NAME) mwaa_auth_token = 'Bearer ' + mwaa_cli_token['CliToken'] mwaa_webserver_hostname = f'https://{mwaa_cli_token["WebServerHostname"]}/aws_mwaa/cli' # 构造完整的CLI命令:dags trigger <dag_name> --conf '<json参数>' raw_data = f'{mwaa_cli_command} {dag_name} --conf \'{dag_conf}\'' # 发送POST请求触发DAG mwaa_response = requests.post( url=mwaa_webserver_hostname, headers={ 'Authorization': mwaa_auth_token, 'Content-Type': 'text/plain' }, data=raw_data ) # 解析并打印返回结果 mwaa_std_err_message = base64.b64decode(mwaa_response.json()['stderr']).decode('utf8') mwaa_std_out_message = base64.b64decode(mwaa_response.json()['stdout']).decode('utf8') print(f"请求状态码: {mwaa_response.status_code}") print('错误信息:\n' + mwaa_std_err_message) print('输出信息:\n' + mwaa_std_out_message)
关键修正说明
原代码中使用params={"x": 11}是错误的,这个参数是给HTTP请求本身的查询参数,并不会传递给DAG。正确的方式是在CLI命令里添加--conf参数,把JSON格式的参数嵌入到命令字符串中。
3. 验证结果
运行触发代码后,登录MWAA的Airflow UI,找到test_param_dag的最新运行实例,进入print_parameter任务的日志页面,就能看到打印的内容:
传入的参数x的值是: 11
内容的提问来源于stack exchange,提问作者Amit Nahmias
相关产品推荐
相关产品推荐

