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

如何正确使用TemplatesServiceClient运行Google Cloud Dataflow模板任务?

问题排查:Dataflow模板任务启动后失败

背景与问题

刚接触Google Cloud,对Dataflow完全陌生。开发了一个Java应用专门用来运行Dataflow任务,代码实现如下,但任务启动后总是失败,错误日志仅显示"Error processing pipeline."。已确认服务账号和工作负载身份配置正确,希望排查代码问题,同时需要TemplatesServiceClient的代码示例来确认是否是权限问题。使用的模板为gs://dataflow-templates/2023-04-11-00_RC00/Cloud_Bigtable_to_GCS_Parquet。

代码实现

配置Bean定义

@Bean(name="dataflowTemplateJobConfigMap")
public Map<DataflowTemplateJob, DataflowTemplateConfig> dataflowTemplateJobConfigMap(
        @Value("${job.gcs.export.gcs.path}") String gcsExportGcsPath,
        @Value("${job.gcs.export.bigtable.project.id}") String gcsExportBigtableProjectId,
        @Value("${job.gcs.export.bigtable.instance.id}") String gcsExportBigtableInstanceId,
        @Value("${job.gcs.export.bigtable.table.id}") String gcsExportBigtableTableId,
        @Value("${job.gcs.export.output.directory}") String gcsExportOutputDirectory,
        @Value("${job.gcs.export.filename.prefix}") String gcsExportFilenamePrefix,
        @Value("${job.gcs.export.num.shards}") String gcsExportNumShards) {
    Map<DataflowTemplateJob, DataflowTemplateConfig> configMap = new HashMap<>();

    // gcs-export job
    DataflowTemplateConfig gcsExportConfig = new DataflowTemplateConfig();
    gcsExportConfig.setGcsPath(gcsExportGcsPath);
    gcsExportConfig.setParameterMap(ImmutableMap.<String, String>builder()
            .put("bigtableProjectId", gcsExportBigtableProjectId)
            .put("bigtableInstanceId", gcsExportBigtableInstanceId)
            .put("bigtableTableId", gcsExportBigtableTableId)
            .put("outputDirectory", gcsExportOutputDirectory)
            .put("filenamePrefix", gcsExportFilenamePrefix)
            .put("numShards", gcsExportNumShards)
            .build());
    configMap.put("GCS_EXPORT", gcsExportConfig);

    return configMap;
}

DataflowTemplateConfig类

@Getter
@Setter
public class DataflowTemplateConfig {
    private String gcsPath;
    private Map<String, String> parameterMap;
}

任务运行方法

@Value("${google.cloud.project.id}")
private String projectId;

@Value("${google.cloud.region}")
private String region;

@Resource(name="dataflowTemplateJobConfigMap")
private Map<DataflowTemplateJob, DataflowTemplateConfig> dataflowTemplateConfigMap;

public Job run(DataflowTemplateJob dataflowTemplateJob) throws IOException {
    DataflowTemplateConfig config = dataflowTemplateConfigMap.get(dataflowTemplateJob);
    try (TemplatesServiceClient templatesServiceClient =TemplatesServiceClient.create()) {
        CreateJobFromTemplateRequest request =
                CreateJobFromTemplateRequest.newBuilder()
                        .setProjectId(projectId)
                        .setJobName(dataflowTemplateJob.getJobName())
                        .putAllParameters(config.getParameterMap())
                        .setEnvironment(RuntimeEnvironment.newBuilder().build())
                        .setLocation(region)
                        .setGcsPath(config.getGcsPath())
                        .build();
        Job response = templatesServiceClient.createJobFromTemplate(request);
        jobRepository.addJob(response);
        return response;
    } catch (Exception e) {
        throw new RuntimeException("Failed to start Dataflow job", e);
    }
}

错误日志

{
insertId: "6e9nqud4qhu"
labels: {4}
logName: "projects/test/test-dataflow/logs/dataflow.googleapis.com%2Fjob-message"
receiveTimestamp: "2023-05-06T00:56:00.125700767Z"
resource: {2}
severity: "ERROR"
textPayload: "Error processing pipeline."
timestamp: "2023-05-06T00:55:59.103993244Z"
}

排查建议与代码示例

代码层面检查点

  • 参数匹配验证:确认Cloud_Bigtable_to_GCS_Parquet模板所需参数名称是否与代码传入的完全一致(部分模板可能使用下划线而非驼峰命名)。
  • RuntimeEnvironment配置:当前代码中RuntimeEnvironment为空白构建,可尝试添加基础配置,例如指定服务账号或机器类型:
    RuntimeEnvironment.newBuilder()
            .setServiceAccountEmail("your-service-account@your-project.iam.gserviceaccount.com")
            .setMachineType("n1-standard-4")
            .build()
    
  • 模板路径有效性:通过gsutil ls命令验证配置的gcsPath是否指向有效模板文件。

TemplatesServiceClient完整示例

import com.google.cloud.dataflow.v1.TemplatesServiceClient;
import com.google.dataflow.v1.CreateJobFromTemplateRequest;
import com.google.dataflow.v1.Job;
import com.google.dataflow.v1.RuntimeEnvironment;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;

public class DataflowTemplateRunner {
    public static void main(String[] args) throws IOException {
        String projectId = "your-project-id";
        String region = "us-central1";
        String templatePath = "gs://dataflow-templates/2023-04-11-00_RC00/Cloud_Bigtable_to_GCS_Parquet";
        String jobName = "bigtable-to-gcs-parquet-job-" + System.currentTimeMillis();

        // 构建模板参数
        Map<String, String> parameters = new HashMap<>();
        parameters.put("bigtableProjectId", "your-bigtable-project");
        parameters.put("bigtableInstanceId", "your-bigtable-instance");
        parameters.put("bigtableTableId", "your-bigtable-table");
        parameters.put("outputDirectory", "gs://your-bucket/output/");
        parameters.put("filenamePrefix", "export-");
        parameters.put("numShards", "5");

        // 创建TemplatesServiceClient
        try (TemplatesServiceClient client = TemplatesServiceClient.create()) {
            RuntimeEnvironment environment = RuntimeEnvironment.newBuilder()
                    .setServiceAccountEmail("your-service-account@your-project.iam.gserviceaccount.com")
                    .build();

            CreateJobFromTemplateRequest request = CreateJobFromTemplateRequest.newBuilder()
                    .setProjectId(projectId)
                    .setLocation(region)
                    .setJobName(jobName)
                    .setGcsPath(templatePath)
                    .putAllParameters(parameters)
                    .setEnvironment(environment)
                    .build();

            Job job = client.createJobFromTemplate(request);
            System.out.printf("Dataflow job started: %s%n", job.getName());
        } catch (Exception e) {
            System.err.printf("Error starting job: %s%n", e.getMessage());
            e.printStackTrace();
        }
    }
}

权限验证步骤

直接使用gcloud命令行启动任务,验证是否能成功运行:

gcloud dataflow jobs run JOB_NAME \
  --gcs-location=gs://dataflow-templates/2023-04-11-00_RC00/Cloud_Bigtable_to_GCS_Parquet \
  --region=REGION \
  --parameters bigtableProjectId=PROJECT_ID,bigtableInstanceId=INSTANCE_ID,bigtableTableId=TABLE_ID,outputDirectory=GS_PATH,filenamePrefix=PREFIX,numShards=SHARDS

如果命令行能成功运行,说明问题大概率出在Java代码的配置或参数传递上;如果命令行也失败,可在Cloud Logging中筛选dataflow.googleapis.com/job-message查看ERROR级别的完整日志。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 00:14:58