因Cloud Composer昂贵,如何用Google Cloud Workflows触发Dataproc Batch任务?
可以用Cloud Workflows编排Dataproc Batch任务,实现依赖控制
完全可以用Google Cloud Workflows替代Cloud Composer来编排带依赖的Dataproc Batch任务,而且能精准控制任务的执行顺序(前一个任务结束后再启动下一个)。下面是具体实现方案:
1. 配置必要权限
确保Workflows使用的服务账号(默认是[PROJECT_NUMBER]@cloudservices.gserviceaccount.com)拥有以下权限:
dataproc.batches.create:创建Dataproc Batch任务dataproc.batches.get:查询Batch任务状态storage.objects.get/storage.objects.list:如果Batch任务依赖GCS中的资源(比如作业文件),需要对应GCS权限- 可通过IAM角色快速配置,比如给服务账号添加
roles/dataproc.editor和roles/storage.objectViewer(根据实际需求调整)
2. Workflows YAML示例(实现两个依赖的Dataproc Batch任务)
下面的YAML定义了一个工作流:先启动第一个Dataproc Batch,等待它执行完成后,再启动第二个Batch任务。
main: steps: # 步骤1:创建第一个Dataproc Batch任务 - create_batch_1: call: http.post args: url: ${"https://dataproc.googleapis.com/v1/projects/" + project + "/regions/" + region + "/batches"} body: job: placement: clusterName: ${cluster_name} pysparkJob: mainPythonFileUri: "gs://your-bucket/path/to/first-job.py" args: ["--input", "gs://input-bucket/data1", "--output", "gs://output-bucket/result1"] batch: batchId: "batch-job-1-${text.replace(text.replace(current_datetime(), ":", "-"), " ", "-")}" headers: Content-Type: "application/json" auth: type: OAuth2 result: batch_1_response # 步骤2:等待第一个Batch任务完成 - wait_for_batch_1: call: http.get args: url: ${batch_1_response.name} auth: type: OAuth2 result: batch_1_status retry: predicate: ${batch_1_status.status.state != "SUCCEEDED" && batch_1_status.status.state != "FAILED"} max_retries: 60 backoff: initial_delay: 30 multiplier: 1.2 max_delay: 120 # 判断第一个任务是否成功,失败则终止工作流 - check_batch_1_success: switch: - condition: ${batch_1_status.status.state == "FAILED"} raise: "第一个Dataproc Batch任务执行失败" # 步骤3:创建第二个Dataproc Batch任务(依赖第一个任务完成) - create_batch_2: call: http.post args: url: ${"https://dataproc.googleapis.com/v1/projects/" + project + "/regions/" + region + "/batches"} body: job: placement: clusterName: ${cluster_name} pysparkJob: mainPythonFileUri: "gs://your-bucket/path/to/second-job.py" args: ["--input", "gs://output-bucket/result1", "--output", "gs://output-bucket/final-result"] batch: batchId: "batch-job-2-${text.replace(text.replace(current_datetime(), ":", "-"), " ", "-")}" headers: Content-Type: "application/json" auth: type: OAuth2 result: batch_2_response # 步骤4:等待第二个Batch任务完成 - wait_for_batch_2: call: http.get args: url: ${batch_2_response.name} auth: type: OAuth2 result: batch_2_status retry: predicate: ${batch_2_status.status.state != "SUCCEEDED" && batch_2_status.status.state != "FAILED"} max_retries: 60 backoff: initial_delay: 30 multiplier: 1.2 max_delay: 120 # 判断第二个任务是否成功 - check_batch_2_success: switch: - condition: ${batch_2_status.status.state == "FAILED"} raise: "第二个Dataproc Batch任务执行失败" # 工作流完成 - finish: return: "所有Dataproc Batch任务执行成功"
3. 变量说明
部署Workflow时可传入以下变量(不建议硬编码在YAML中):
project:你的GCP项目IDregion:Dataproc集群所在区域(比如us-central1)cluster_name:要运行Batch任务的Dataproc集群名称
4. 部署与触发Workflow
- 部署命令:
gcloud workflows deploy dataproc-batch-workflow \ --source=workflow.yaml \ --region=us-central1 \ --set-env-vars=project=your-project-id,region=us-central1,cluster_name=your-cluster-name - 触发Workflow:
gcloud workflows run dataproc-batch-workflow --region=us-central1
关键逻辑说明
- 任务依赖控制:通过
wait_for_batch_1步骤循环查询第一个Batch的状态,直到它进入SUCCEEDED或FAILED状态,之后才会执行create_batch_2步骤,完全满足“一个任务结束后再启动另一个”的需求。 - 重试机制:使用Workflows的
retry配置,每隔30秒(最多延迟到120秒)查询一次任务状态,最多重试60次(对应最长等待3600秒,可根据任务时长调整)。 - 错误处理:如果任一Batch任务失败,工作流会抛出异常终止,便于排查问题。
内容的提问来源于stack exchange,提问作者no-stale-reads
相关产品推荐
相关产品推荐

