是否可以通过Cloud Composer的Python客户端触发DAG?
是否可以通过Cloud Composer Python客户端触发DAG
可以实现,但你查的Composer管理API客户端本身没有直接封装对应方法,原因和实现方案如下:
核心原因
gcloud composer environments run命令不是直接调用Composer的原生管理接口,而是gcloud封装的组合操作:执行该命令时gcloud会先调用Composer管理API拉取对应环境的Airflow Web Server地址、身份验证配置,再调用Airflow本身的REST API完成DAG触发,因此你在Composer管理API的文档里找不到直接对应的方法。
两种实现方案
- 方案1:组合使用Composer管理客户端 + Airflow REST API自行实现
- 用你正在使用的Composer Python客户端查询目标环境的Airflow Web服务器端点
- 按照你使用的Composer版本对应的身份验证规则,构造Airflow REST API请求,调用POST
/api/v1/dags/{dag_id}/dagRuns接口即可触发DAG,请求参数里可以指定run_id等配置
- 方案2:直接使用Airflow服务对应的Python客户端调用触发接口
可以直接调用Google官方封装的Airflow服务Python客户端,直接传入环境信息、DAG ID、运行参数即可完成触发,不需要自己拼接HTTP请求。
简化示例(适用于Composer 2环境)
from google.cloud import composer_v1 import requests from google.oauth2 import service_account from google.auth.transport.requests import AuthorizedSession # 初始化Composer客户端查询环境信息 composer_client = composer_v1.EnvironmentsClient() environment_path = composer_client.environment_path("你的项目ID", "环境所在区域", "环境名称") environment = composer_client.get_environment(name=environment_path) airflow_endpoint = environment.config.airflow_uri # 构造授权会话 credentials, _ = service_account.default(scopes=["https://www.googleapis.com/auth/cloud-platform"]) authed_session = AuthorizedSession(credentials) # 调用Airflow API触发DAG dag_id = "你要触发的DAG ID" trigger_url = f"{airflow_endpoint}/api/v1/dags/{dag_id}/dagRuns" payload = { "run_id": "自定义的run_id", "conf": {} # 此处填入需要传给DAG的运行参数 } response = authed_session.post(trigger_url, json=payload) print(response.json())
内容的提问来源于stack exchange,提问作者jamiet
相关产品推荐
相关产品推荐

