如何在Cloud Function触发的Composer DAG中使用传入的data参数?
在Composer DAG中获取Cloud Function触发时传入的data参数
1. 确认Cloud Function中data参数的传递格式
在调用trigger_dag_gcf接口时,需将额外数据封装在请求体的data字段中,示例代码如下:
import google.auth from googleapiclient.discovery import build def trigger_dag(event, context): credentials, project_id = google.auth.default() service = build('composer', 'v1', credentials=credentials) dag_id = "your-target-dag-id" env_full_path = "projects/[PROJECT_ID]/locations/[REGION]/environments/[ENV_NAME]" request_body = { "data": { "user_id": "12345", "file_path": "/bucket/path/file.csv" } } request = service.projects().locations().environments().dags().trigger( name=f"{env_full_path}/dags/{dag_id}", body=request_body ) response = request.execute() return response
2. 在DAG中读取传入的数据
Cloud Function传入的data会被自动封装到Airflow的dag_run.conf对象中,你可以通过两种方式获取:
方式一:在PythonOperator中通过上下文读取
给PythonOperator开启provide_context=True后,就能从上下文参数中获取dag_run.conf:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def handle_trigger_data(**context): # 提取传入的参数 trigger_conf = context["dag_run"].conf user_id = trigger_conf.get("user_id") file_path = trigger_conf.get("file_path") # 执行业务逻辑 print(f"Processing file for user {user_id}: {file_path}") with DAG( dag_id="your-target-dag-id", start_date=datetime(2024, 1, 1), schedule_interval=None, catchup=False ) as dag: process_task = PythonOperator( task_id="process_trigger_data", python_callable=handle_trigger_data, provide_context=True )
方式二:用Jinja模板在其他Operator中引用
如果是BashOperator这类支持Jinja模板的Operator,可以直接通过模板语法调用dag_run.conf中的字段:
from airflow.operators.bash import BashOperator bash_task = BashOperator( task_id="echo_trigger_data", bash_command='echo "User ID: {{ dag_run.conf.user_id }}, File Path: {{ dag_run.conf.file_path }}"' )
注意事项
- 确保Cloud Function使用的服务账号拥有
composer.environments.dags.trigger权限 - 传入的
data需可JSON序列化,否则会触发接口调用失败
内容的提问来源于stack exchange,提问作者le Minh Nguyen
相关产品推荐
相关产品推荐

