Dataflow Pipeline执行失败:无法打开GCS文件的问题求助
问题描述
我在网上找到一个Demo尝试复现,运行以下命令时:
python pipeline.py --streaming --runner DataflowRunner \ --project <PROJECT> \ --temp_location gs://tweeps-stream/temp \ --staging_location gs://tweeps-stream/staging \ --region us-west1 \ --job_name tweeps
出现如下错误:
apache_beam.runners.dataflow.dataflow_runner.DataflowRuntimeException: Dataflow pipeline failed. State: FAILED, Error: Unable to open file: gs://tweeps-stream/staging/tweep.1697803634.304288/pipeline.pb.
我通过GCP云控制台操作,已给自己分配Dataflow、BigQuery和GCS的管理员权限,本地运行python pipeline.py --streaming正常。使用的代码如下:
from apache_beam.options.pipeline_options import PipelineOptions from sys import argv import apache_beam as beam import argparse PROJECT_ID = 'gcp-project' SUBSCRIPTION = 'projects/' + PROJECT_ID + '/subscriptions/tweeps' SCHEMA = 'created_at:TIMESTAMP,tweep_id:STRING,text:STRING,user:STRING,flagged:BOOLEAN' def parse_pubsub(data): # use the json library to convert the datafrom pubsub to a python dictionary object # makes it much easier to manipulate and Apache Beam also supports python dictionaries directly to BigQuery. import json return json.loads(data) def fix_timestamp(data): import datetime d = datetime.datetime.strptime(data['created_at'], "%d/%b/%Y:%H:%M:%S") data['created_at'] = d.strftime("%Y-%m-%d %H:%M:%S") return data def check_tweep(data): BAD_WORDS = ['attack', 'drug', 'gun'] data['flagged'] = False for word in BAD_WORDS: if word in data['text']: data['flagged'] = True return data if __name__ == '__main__': parser = argparse.ArgumentParser() known_args = parser.parse_known_args(argv) p = beam.Pipeline(options=PipelineOptions()) (p | 'ReadData' >> beam.io.ReadFromPubSub(subscription=SUBSCRIPTION).with_output_types(bytes) | 'Decode' >> beam.Map(lambda x: x.decode('utf-8')) | 'PubSubToJSON' >> beam.Map(parse_pubsub) | 'FixTimestamp' >> beam.Map(fix_timestamp) | 'CheckTweep' >> beam.Map(check_tweep) | 'WriteToBigQuery' >> beam.io.WriteToBigQuery( '{0}:tweeper.tweeps'.format(PROJECT_ID), schema=SCHEMA, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)) result = p.run() result.wait_until_finish()
疑问:已配置正确权限,本地运行正常,为何会出现该错误?
解决建议
- 修正PipelineOptions初始化逻辑:代码中
PipelineOptions()未加载命令行参数,导致Dataflow的运行配置(如runner、project、存储路径等)未被正确读取,这是核心问题。将初始化代码改为:
本地运行时默认使用DirectRunner,不需要额外参数所以能正常执行,但提交Dataflow时必须让PipelineOptions读取命令行传入的参数,否则无法正确生成并上传p = beam.Pipeline(options=PipelineOptions(argv[1:]))pipeline.pb文件。 - 检查Dataflow服务账号权限:确认Dataflow服务账号(格式为
service-<PROJECT_NUMBER>@dataflow-service-producer-prod.iam.gserviceaccount.com)对gs://tweeps-stream存储桶拥有Storage Object Admin权限。即使你个人有管理员权限,Dataflow运行时使用的是服务账号,权限不足会导致无法读写staging目录下的文件。 - 验证存储桶与参数正确性:确保
gs://tweeps-stream存储桶已创建,且命令中的<PROJECT>已替换为实际的GCP项目ID,避免路径解析错误。 - 清理旧缓存文件:删除
gs://tweeps-stream/temp和gs://tweeps-stream/staging下的所有旧文件,重新运行命令,避免残留的旧文件导致冲突。 - 对齐Beam版本:确保本地使用的Apache Beam版本与Dataflow支持的版本一致,版本不匹配可能导致序列化异常,无法生成有效的
pipeline.pb文件。
内容的提问来源于stack exchange,提问作者Mike3355
相关产品推荐
相关产品推荐

