如何用Java代码手动触发Cloud Composer DAG运行
用Java本地触发Cloud Composer中手动DAG的运行
前提准备
- 本地已通过
gcloud auth application-default login配置好应用默认凭据 - 掌握你的Cloud Composer环境的Airflow Web UI端点(格式类似
https://<your-composer-environment-web-ui-url>) - 确认目标DAG的ID(即DAG代码中定义的
dag_id)
依赖配置(Maven)
在pom.xml中添加以下依赖,用于处理凭据、HTTP请求和JSON解析:
<dependencies> <!-- Google Cloud 应用默认凭据处理 --> <dependency> <groupId>com.google.auth</groupId> <artifactId>google-auth-library-oauth2-http</artifactId> <version>1.20.0</version> </dependency> <!-- HTTP客户端 --> <dependency> <groupId>com.google.http-client</groupId> <artifactId>google-http-client</artifactId> <version>1.43.3</version> </dependency> <!-- JSON工具 --> <dependency> <groupId>com.google.code.gson</groupId> <artifactId>gson</artifactId> <version>2.10.1</version> </dependency> </dependencies>
核心代码实现
由于Google官方Java客户端库google-cloud-orchestration-airflow仅封装了Composer环境管理能力,未覆盖Airflow的DAG操作API,因此直接调用Airflow REST API是最直接的方案:
import com.google.auth.oauth2.GoogleCredentials; import com.google.auth.oauth2.IdTokenProvider; import com.google.gson.JsonObject; import com.google.gson.JsonParser; import com.google.api.client.http.*; import com.google.api.client.http.javanet.NetHttpTransport; import java.io.IOException; import java.util.Collections; public class ComposerDagTrigger { public static void main(String[] args) throws IOException { // 替换为你的实际配置 String airflowWebUiUrl = "https://your-composer-environment-web-ui-url"; String dagId = "your-target-dag-id"; // 1. 获取应用默认凭据并生成Airflow端点的ID Token GoogleCredentials credentials = GoogleCredentials.getApplicationDefault() .createScoped(Collections.singletonList("openid")); IdTokenProvider idTokenProvider = (IdTokenProvider) credentials; String idToken = idTokenProvider.createIdToken(airflowWebUiUrl, null).getTokenValue(); // 2. 构建触发请求体,可自定义运行ID和参数 JsonObject requestBody = new JsonObject(); requestBody.addProperty("dag_run_id", "manual-trigger-" + System.currentTimeMillis()); // 可选:添加DAG运行参数 // JsonObject conf = new JsonObject(); // conf.addProperty("param_key", "param_value"); // requestBody.add("conf", conf); // 3. 构建并发送POST请求 HttpTransport httpTransport = new NetHttpTransport(); HttpRequestFactory requestFactory = httpTransport.createRequestFactory(); GenericUrl triggerUrl = new GenericUrl(airflowWebUiUrl + "/api/v1/dags/" + dagId + "/dagRuns"); HttpRequest request = requestFactory.buildPostRequest(triggerUrl, ByteArrayContent.fromString( "application/json", requestBody.toString())); request.getHeaders().setAuthorization("Bearer " + idToken); // 4. 处理响应 HttpResponse response = request.execute(); try { int statusCode = response.getStatusCode(); if (statusCode >= 200 && statusCode < 300) { JsonObject jsonResponse = JsonParser.parseString(response.parseAsString()).getAsJsonObject(); System.out.println("DAG触发成功,运行ID:" + jsonResponse.get("dag_run_id").getAsString()); } else { System.err.println("触发失败,状态码:" + statusCode + ",响应内容:" + response.parseAsString()); } } finally { response.disconnect(); } } }
关键说明
- 凭据逻辑和Python一致:通过应用默认凭据生成ID Token,无需额外配置密钥文件
- 运行ID建议用唯一值(如时间戳),避免重复触发的冲突
- 若DAG需要接收参数,在
conf字段中添加键值对即可,对应Airflow运行时的dag_run.conf
内容的提问来源于stack exchange,提问作者Sasirekha MSVL
相关产品推荐
相关产品推荐

