DataFlowRunner运行Dataflow管道报错,DirectRunner正常求助
问题
使用Dataflow时出现Error processing pipeline报错,DirectRunner运行正常,切换为DataFlowRunner时失败。错误日志如下:
ERROR:apache_beam.runners.dataflow.dataflow_runner:Console URL: https://console.cloud.google.com/dataflow/jobs/<RegionId>/2023-04-05_02_00_47-7238238223888513941?project=<ProjectId> Traceback (most recent call last): File "test.py", line 51, in <module> run() File "test.py", line 44, in run output | 'Write' >> WriteToText("gs://<bucket>/output/wc.txt") File "/opt/py38/lib64/python3.8/site-packages/apache_beam/pipeline.py", line 601, in __exit__ self.result.wait_until_finish() File "/opt/py38/lib64/python3.8/site-packages/apache_beam/runners/dataflow/dataflow_runner.py", line 1555, in wait_until_finish raise DataflowRuntimeException( apache_beam.runners.dataflow.dataflow_runner.DataflowRuntimeException: Dataflow pipeline failed. State: FAILED, Error: Error processing pipeline.
当前已配置的IAM权限:
- Storage Object Administrator
- BigqueryConnectionCustom Role
- SQL Cloud Clients
- BigQuery data editor
- Dataflow developer
- Service account user
- BigQuery user
- Compute Viewer
- Worker Dataflow
运行代码:
import os import apache_beam as beam from apache_beam.io import WriteToText from apache_beam.options.pipeline_options import PipelineOptions def pipelineOptions(pipeline_args): pipeline_options = PipelineOptions( pipeline_args, runner="DirectRunner", project=<project-name>, job_name="testbigquery", temp_location=<temp-location>, region=<region> ) return pipeline_options def run(argv=None): print("Start Process") pipeline_options = pipelineOptions(argv) pipeline = beam.Pipeline(options=pipeline_options) with pipeline as p: lines = p counts = ( lines | 'Split' >> (beam.Create(["test", "fix", "test"])) | 'PairWithOne' >> beam.Map(lambda x: (x, 1)) | 'GroupAndSum' >> beam.CombinePerKey(sum)) def format_result(word, count): return '%s: %d' % (word, count) output = counts | 'Format' >> beam.MapTuple(format_result) output | 'Write' >> WriteToText("gs://<bucket>/output/wc.txt") print("End Process") if __name__ == '__main__': run()
故障排查与解决思路
修正Runner配置
代码中pipelineOptions函数硬编码了runner="DirectRunner",切换到DataFlowRunner时需修改此处为"DataflowRunner",或通过命令行参数覆盖配置(确保命令行参数优先级高于代码硬设置)。验证存储路径配置
- 确认
temp_location指向的GCS桶存在,且运行Dataflow的账号拥有该桶的读写权限,同时桶区域需与Dataflow区域一致,避免跨区域访问问题。 - 修改
WriteToText的输出路径:Dataflow要求输出路径为目录而非具体文件名,建议改为gs://<bucket>/output/,WriteToText会自动生成带后缀的输出文件。
- 确认
检查Dataflow服务账号权限
- 确认Dataflow默认服务账号(格式为
service-<project-number>@dataflow-service-producer-prod.iam.gserviceaccount.com)拥有以下权限:- Dataflow Worker角色
- 目标GCS桶的Storage Object Creator/Editor权限
- 确保项目已启用Dataflow API。
- 确认Dataflow默认服务账号(格式为
查看详细错误日志
现有错误信息过于笼统,需登录GCP控制台进入对应Dataflow作业页面,查看Worker日志、作业详情中的错误事件,定位具体失败原因(如依赖缺失、资源不足、网络配置问题等)。确认版本兼容性
检查使用的Apache Beam版本是否与Dataflow服务兼容,建议使用官方推荐的稳定版本,避免版本不匹配导致的运行时异常。
内容的提问来源于stack exchange,提问作者luca
相关产品推荐
相关产品推荐

