You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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、存储路径等)未被正确读取,这是核心问题。将初始化代码改为:
    p = beam.Pipeline(options=PipelineOptions(argv[1:]))
    
    本地运行时默认使用DirectRunner,不需要额外参数所以能正常执行,但提交Dataflow时必须让PipelineOptions读取命令行传入的参数,否则无法正确生成并上传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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 04:43:24