Dataflow Prime作业在Transform上配置资源提示后运行失败
问题描述
我使用Apache Beam Java SDK v2.32.0编写了一个经典Dataflow模板,该模板的功能是从Pub/Sub订阅消费消息并写入Google Cloud Storage。
通过--additional-experiments enable_prime启用Dataflow Prime实验特性,同时通过--parameters=resourceHints=min_ram=8GiB配置流水线级资源提示时,模板可以正常运行作业,运行命令如下:
gcloud dataflow jobs run my-job-name \ --additional-experiments enable_prime \ --disable-public-ips \ --gcs-location gs://bucket/path/to/template \ --num-workers 1 \ --max-workers 16 \ --parameters=resourceHints=min_ram=8GiB,other_pipeline_options=true \ --project my-project \ --region us-central1 \ --service-account-email my-service-account@my-project.iam.gserviceaccount.com \ --staging-location gs://bucket/path/to/staging \ --subnetwork https://www.googleapis.com/compute/v1/projects/my-project/regions/us-central1/subnetworks/my-subnet
为了使用Dataflow Prime的Right Fitting能力,我修改了流水线代码,在FileIO Transform上添加了资源提示,修改后的代码如下:
class WriteGcsFileTransform extends PTransform<PCollection<Input>, WriteFilesResult<Destination>> { private static final long serialVersionUID = 1L; @Override public WriteFilesResult<Destination> expand(PCollection<Input> input) { return input.apply( FileIO.<Destination, Input>writeDynamic() .by(myDynamicDestinationFunction) .withDestinationCoder(Destination.coder()) .withNumShards(8) .withNaming(myDestinationFileNamingFunction) .withTempDirectory("gs://bucket/path/to/temp") .withCompression(Compression.GZIP) .setResourceHints(ResourceHints.create().withMinRam("32GiB")) ); }
基于修改后代码生成的模板运行作业时,作业持续进入崩溃循环无法正常启动,重复出现的错误日志如下:
{ "insertId": "s=97e1ecd30e0243609d555685318325b4;i=4e1;b=6c7f5d65f3994eada5f20672dab1daf1;m=912f16c;t=5d024689cb030;x=b36751718b3d80c1", "jsonPayload": { "line": "pod_workers.go:191", "message": "Error syncing pod 4cf7cbf98df4b5e2d054abce7da1262b (\"df-df-hvm-my-job-name-11061310-qn51-harness-jb9f_default(4cf7c6bf982df4b5eb2d054abce7da12)\"), skipping: failed to \"StartContainer\" for \"artifact\" with CrashLoopBackOff: \"back-off 40s restarting failed container=artifact pod=df-df-hvm-my-job-name-11061310-qn51-harness-jb9f_default(4cf7c6bf982df4b5eb2d054abce7da12)\"", "thread": "807" }, "resource": { "type": "dataflow_step", "labels": { "project_id": "my-project", "region": "us-central1", "step_id": "", "job_id": "2021-11-06_12_10_27-510057810808146686", "job_name": "my-job-name" } }, "timestamp": "2021-11-06T20:14:36.052491Z", "severity": "ERROR", "labels": { "compute.googleapis.com/resource_type": "instance", "dataflow.googleapis.com/log_type": "system", "compute.googleapis.com/resource_id": "4695846446965678007", "dataflow.googleapis.com/job_name": "my-job-name", "dataflow.googleapis.com/job_id": "2021-11-06_12_10_27-510057810808146686", "dataflow.googleapis.com/region": "us-central1", "dataflow.googleapis.com/service_option": "prime", "compute.googleapis.com/resource_name": "df-hvm-my-job-name-11061310-qn51-harness-jb9f" }, "logName": "projects/my-project/logs/dataflow.googleapis.com%2Fkubelet", "receiveTimestamp": "2021-11-06T20:14:46.471285909Z" }
请问我在Transform上使用资源提示的方式是否存在错误?
解答
你调用setResourceHints的语法本身符合Apache Beam的API规范,报错是版本兼容和配额限制两个问题共同导致的:
- 你使用的Apache Beam Java SDK v2.32.0对Transform层级的资源提示支持存在缺陷,只有全局级别的
resourceHints参数可以被Dataflow Prime正确解析,Transform级别的资源提示会在调度阶段生成无效的Pod资源配置,直接导致harness容器启动失败。 - 你配置的32GiB单Transform最小内存,超过了us-central1区域Dataflow Prime默认支持的单vCPU绑定内存上限,未提前调整配额的情况下,调度器无法分配符合要求的计算资源,就会触发容器循环崩溃。
修复方案
- 升级Apache Beam Java SDK到v2.37.0及以上版本,该版本修复了Transform层级资源提示和Dataflow Prime的适配问题。
- 若暂时无法升级SDK,可将Transform级的内存要求调整为全局资源提示,配合调整对应步骤的并行度实现同等效果。
- 确认项目在对应区域的Dataflow Prime内存配额满足32GiB的要求,配额不足可在GCP控制台配额页面申请上调。
内容的提问来源于stack exchange,提问作者Brent Worden
相关产品推荐
相关产品推荐

