You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 23:02:23