如何通过Java应用基于GCS模板启动GCP Dataflow作业
问题场景
- 现有Java编写的Apache Beam Dataflow作业,当前在GCP上可通过以下流程正常运行:
- 从业务代码生成Dataflow模板,上传至Cloud Storage存储桶
- 在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实例。
解决方案
你的整体实现思路是正确的,只需要修正客户端构造逻辑、补充缺失依赖、完善调用参数即可完成需求,具体步骤如下:
- 补充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>
- 修正Dataflow客户端构造逻辑
替换错误传入的LowLevelHttpRequest实例,使用合法的HttpTransport实现构造客户端:
// 初始化可信HTTP传输实例 HttpTransport httpTransport = ApacheHttpTransport.newTrustedTransport(); // 构造Dataflow服务客户端 Dataflow dataflow = new Dataflow.Builder(httpTransport, JacksonFactory.getDefaultInstance(), credential) .setApplicationName("my-job") .build();
- 完善模板启动调用逻辑
你当前的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
相关产品推荐
相关产品推荐

