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

如何通过Java应用基于GCS模板启动GCP Dataflow作业

问题场景
  • 现有Java编写的Apache Beam Dataflow作业,当前在GCP上可通过以下流程正常运行:
    1. 从业务代码生成Dataflow模板,上传至Cloud Storage存储桶
    2. 在GCP控制台的Dataflow->Jobs页面,选择从模板创建作业的选项直接启动任务
  • 目标需求:通过Java应用对外提供API接口,接口收到请求时,自动基于已存储在Cloud Storage中的Dataflow模板启动作业。
  • 已知可通过REST API POST /v1b3/projects/project_id/locations/loc/templates:launch?gcsPath=template-location 实现需求,但未找到对应的Java参考实现,自行编码时遇到运行报错。
问题复现

引入依赖

首先在Spring Boot项目中引入了如下Maven依赖:

<!-- https://mvnrepository.com/artifact/com.google.apis/google-api-services-dataflow -->
<dependency>
    <groupId>com.google.apis</groupId>
    <artifactId>google-api-services-dataflow</artifactId>
    <version>v1b3-rev20210825-1.32.1</version>
</dependency>

核心实现代码

在Controller层编写的createJob方法如下:

public static void createJob() throws IOException {
    GoogleCredential credential = GoogleCredential.fromStream(new FileInputStream("myCertKey.json")).createScoped(
        java.util.Arrays.asList("https://www.googleapis.com/auth/cloud-platform"));       
    try{
        // 该行抛出错误
        Dataflow dataflow = new Dataflow.Builder(new LowLevelHttpRequest(), new JacksonFactory(),
        credential).setApplicationName("my-job").build();           
        
        //RuntimeEnvironment
        RuntimeEnvironment env = new RuntimeEnvironment();
        env.setBypassTempDirValidation(false);
        // 已添加所有环境配置项

        //parameters
        HashMap<String,String> params = new HashMap<>();
        params.put("bigtableEmulatorPort", "-1"); 
        params.put("gcsPath", "gs://bucket//my.json");
        // 已添加所有其他作业参数
        
        LaunchTemplateParameters content = new LaunchTemplateParameters();
        content.setJobName("Test-job");
        content.setEnvironment(env);
        content.setParameters(params);
        
        dataflow.projects().locations().templates().launch("project-id", "location", content);           
    }catch (Exception e){
        log.info("error occured", e);
    }
}

报错信息

代码运行到构建Dataflow实例的行时抛出如下错误:

{"id":null,"message":"'boolean com.google.api.client.http.HttpTransport.isMtls()'"}

已初步排查到报错直接原因是Dataflow.Builder构造方法第一个入参要求为HttpTransport类型,代码中错误传入了LowLevelHttpRequest实例。

解决方案

你的整体实现思路是正确的,只需要修正客户端构造逻辑、补充缺失依赖、完善调用参数即可完成需求,具体步骤如下:

  1. 补充HTTP传输实现依赖
    你当前仅引入了Dataflow API客户端依赖,缺少Google HTTP客户端对应的传输层实现,需要在pom.xml中补充以下依赖:
<dependency>
    <groupId>com.google.http-client</groupId>
    <artifactId>google-http-client-apache-v2</artifactId>
    <!-- 版本需与dataflow依赖传递引入的google-http-client核心版本保持一致,避免版本不匹配触发isMtls()方法找不到的错误 -->
    <version>1.42.3</version>
</dependency>
  1. 修正Dataflow客户端构造逻辑
    替换错误传入的LowLevelHttpRequest实例,使用合法的HttpTransport实现构造客户端:
// 初始化可信HTTP传输实例
HttpTransport httpTransport = ApacheHttpTransport.newTrustedTransport();
// 构造Dataflow服务客户端
Dataflow dataflow = new Dataflow.Builder(httpTransport, JacksonFactory.getDefaultInstance(), credential)
        .setApplicationName("my-job")
        .build();
  1. 完善模板启动调用逻辑
    你当前的launch调用没有传入GCS模板路径,也没有执行实际的请求发送,需要补充参数并调用execute()方法触发请求:
// 替换为实际的项目ID、区域、GCS模板路径
String projectId = "your-gcp-project-id";
String region = "dataflow-job-region";
String gcsTemplatePath = "gs://your-bucket/path/to/dataflow-template";

// 组装启动请求
Dataflow.Projects.Locations.Templates.Launch launchRequest = dataflow.projects()
        .locations()
        .templates()
        .launch(projectId, region, content)
        .setGcsPath(gcsTemplatePath);

// 发送请求,获取作业启动结果
LaunchTemplateResponse response = launchRequest.execute();
// 可从返回结果中获取启动后的作业ID等信息
String launchedJobId = response.getJob().getId();

注意事项

  • 密钥文件对应的服务账号需要拥有目标项目的Dataflow作业启动权限(最小权限为roles/dataflow.developer,也可使用roles/dataflow.admin),同时需要拥有GCS模板文件的读取权限,以及作业运行涉及的其他GCP资源(Bigtable、GCS、BigQuery等)的对应访问权限。
  • Dataflow作业名需符合命名规范:仅支持小写字母、数字、连字符,长度1-63位,不能以连字符开头或结尾,同一区域下作业名需唯一。
  • 如果Spring Boot应用运行在GCP环境内(GCE、GKE、Cloud Run、App Engine等),不需要手动加载本地密钥文件,可直接使用环境默认服务账号凭证,安全性更高。

内容的提问来源于stack exchange,提问作者Sweety

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 20:36:29