为何我的DoFn会无限重启元素处理过程?
问题:Dataflow中长时运行网页抓取DoFn无限重启无法推进流水线
当处理运行时间2-3小时的大型元素时,网页抓取DoFn返回新pCollection后会无限重启执行,无法进入流水线下一步。
错误信息
主要错误
"Commit failed, will be retried at higher level but may not succeed" computation = 'P51', sharding_key = '12c92e1638acbb09', status = generic::internal: Windmill failed to commit the work item. CommitStatus: VALIDATION_FAILED"
完整错误日志
[ { "insertId": "7054061089043139090:270578:0:287794", "jsonPayload": { "job": "2022-08-19_13_38_49-15282523611410283510", "work": "2", "step": "Scrape Products", "instruction": "process_bundle-2-1", "thread": "Thread-14", "logger": "REDACTED PATH/product_scraper.py:69", "worker": "df-hvm-beamapp-vincentye-0819203-08191338-dtgt-harness-l4gg", "portability_worker_id": "sdk-0-0", "message": "returning new_products_dataframe and REDACTED_dataframe" }, "resource": { "type": "dataflow_step", "labels": { "job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "step_id": "Scrape Products", "project_id": "REDACTED", "job_id": "2022-08-19_13_38_49-15282523611410283510", "region": "us-west2" } }, "timestamp": "2022-08-19T21:38:39.037309408Z", "severity": "WARNING", "labels": { "dataflow.googleapis.com/region": "us-west2", "dataflow.googleapis.com/job_id": "2022-08-19_13_38_49-15282523611410283510", "compute.googleapis.com/resource_name": "df-hvm-beamapp-vincentye-0819203-08191338-dtgt-harness-l4gg", "dataflow.googleapis.com/job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "dataflow.googleapis.com/service_option": "prime", "compute.googleapis.com/resource_id": "7054061089043139090", "compute.googleapis.com/resource_type": "instance" }, "logName": "projects/REDACTED/logs/dataflow.googleapis.com%2Fworker", "receiveTimestamp": "2022-08-19T21:38:52.093172441Z" }, { "insertId": "s=6f21f5e3f31c4dfca9b6c2dc7f2de30a;i=664;b=f6f95ab5a8234fdeac1c3d31e7f63933;m=d31e6263;t=5e69eebe139ed;x=77b6120c808dd2fb", "jsonPayload": { "line": "sampler.go:311", "message": "Successfully sampled resources" }, "resource": { "type": "dataflow_step", "labels": { "job_id": "2022-08-19_13_38_49-15282523611410283510", "project_id": "REDACTED", "region": "us-west2", "step_id": "", "job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u" } }, "timestamp": "2022-08-19T21:38:40.180661Z", "severity": "INFO", "labels": { "dataflow.googleapis.com/region": "us-west2", "dataflow.googleapis.com/log_type": "system", "dataflow.googleapis.com/job_id": "2022-08-19_13_38_49-15282523611410283510", "compute.googleapis.com/resource_id": "4422558927708899858", "dataflow.googleapis.com/job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "compute.googleapis.com/resource_name": "df-hvm-beamapp-vincentye-0819203-08191338-dtgt-harness-3jp2", "compute.googleapis.com/resource_type": "instance", "dataflow.googleapis.com/service_option": "prime" }, "logName": "projects/REDACTED/logs/dataflow.googleapis.com%2Fresource", "receiveTimestamp": "2022-08-19T21:38:41.688521247Z" }, { "insertId": "7054061089043139090:270576:0:17358", "jsonPayload": { "message": "Commit failed, will be retried at higher level but may not succeed. computation = \"P51\", sharding_key = \"62fa5cb5ab454560\", status = generic::internal: Windmill failed to commit the work item. CommitStatus: VALIDATION_FAILED", "thread": "119", "line": "streaming_worker_client.cc:514" }, "resource": { "type": "dataflow_step", "labels": { "region": "us-west2", "project_id": "REDACTED", "job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "job_id": "2022-08-19_13_38_49-15282523611410283510", "step_id": "" } }, "timestamp": "2022-08-19T21:38:40.686591Z", "severity": "ERROR", "labels": { "dataflow.googleapis.com/service_option": "prime", "dataflow.googleapis.com/job_id": "2022-08-19_13_38_49-15282523611410283510", "dataflow.googleapis.com/log_type": "system", "dataflow.googleapis.com/job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "dataflow.googleapis.com/region": "us-west2", "compute.googleapis.com/resource_name": "df-hvm-beamapp-vincentye-0819203-08191338-dtgt-harness-l4gg", "compute.googleapis.com/resource_type": "instance", "compute.googleapis.com/resource_id": "7054061089043139090" }, "logName": "projects/REDACTED/logs/dataflow.googleapis.com%2Fharness", "receiveTimestamp": "2022-08-19T21:38:42.095203701Z" }, { "insertId": "7054061089043139090:270575:0:13058", "jsonPayload": { "line": "exec.go:66", "message": "E0819 21:38:40.686591 119 streaming_worker_client.cc:514] Commit failed, will be retried at higher level but may not succeed. computation = \"P51\", sharding_key = \"62fa5cb5ab454560\", status = generic::internal: Windmill failed to commit the work item. CommitStatus: VALIDATION_FAILED" }, "resource": { "type": "dataflow_step", "labels": { "job_id": "2022-08-19_13_38_49-15282523611410283510", "step_id": "", "job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "project_id": "REDACTED", "region": "us-west2" } }, "timestamp": "2022-08-19T21:38:40.686846Z", "severity": "INFO", "labels": { "dataflow.googleapis.com/log_type": "system", "dataflow.googleapis.com/service_option": "prime", "compute.googleapis.com/resource_type": "instance", "dataflow.googleapis.com/region": "us-west2", "compute.googleapis.com/resource_name": "df-hvm-beamapp-vincentye-0819203-08191338-dtgt-harness-l4gg", "dataflow.googleapis.com/job_id": "2022-08-19_13_38_49-15282523611410283510", "compute.googleapis.com/resource_id": "7054061089043139090", "dataflow.googleapis.com/job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u" }, "logName": "projects/REDACTED/logs/dataflow.googleapis.com%2Fharness-startup", "receiveTimestamp": "2022-08-19T21:38:52.093608994Z" }, { "insertId": "7054061089043139090:270576:0:17641", "jsonPayload": { "line": "streaming_worker_client.cc:566", "thread": "119", "message": "Error while processing a work item: INTERNAL: Windmill failed to commit the work item. CommitStatus: VALIDATION_FAILED\n=== Source Location Trace: ===\ndist_proc/dax/workflow/worker/streaming/streaming_worker_client.cc:492" }, "resource": { "type": "dataflow_step", "labels": { "step_id": "", "job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "job_id": "2022-08-19_13_38_49-15282523611410283510", "project_id": "REDACTED", "region": "us-west2" } }, "timestamp": "2022-08-19T21:38:40.687126Z", "severity": "WARNING", "labels": { "compute.googleapis.com/resource_type": "instance", "dataflow.googleapis.com/region": "us-west2", "compute.googleapis.com/resource_name": "df-hvm-beamapp-vincentye-0819203-08191338-dtgt-harness-l4gg", "dataflow.googleapis.com/job_name": "beamapp-vincentye-0819203841-963772-8quxtp6u", "dataflow.googleapis.com/job_id": "2022-08-19_13_38_49-15282523611410283510", "compute.googleapis.com/resource_id": "7054061089043139090", "dataflow.googleapis.com/service_option": "prime", "dataflow.googleapis.com/log_type": "system" }, "logName": "projects/REDACTED/logs/dataflow.googleapis.com%2Fharness", "receiveTimestamp": "2022-08-19T21:39:02.095105228Z" } ]
更多信息
- 仅处理运行时间超2小时的大型元素时触发错误
- 此类元素生成的pCollection大小在100-150MB之间,包含以字节形式存储图片的Pandas DataFrame
- 失败元素会被无限重试
- LocalRunner运行流水线时无此问题
已尝试的解决方法
- 关闭Dataflow Prime,将machine_type设置为
n2-highmem-2 - 将
number_of_worker_harness_threads设置为3
内容的提问来源于stack exchange,提问作者Vincent Ye
相关产品推荐
相关产品推荐

