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

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()

故障排查与解决思路

  1. 修正Runner配置
    代码中pipelineOptions函数硬编码了runner="DirectRunner",切换到DataFlowRunner时需修改此处为"DataflowRunner",或通过命令行参数覆盖配置(确保命令行参数优先级高于代码硬设置)。

  2. 验证存储路径配置

    • 确认temp_location指向的GCS桶存在,且运行Dataflow的账号拥有该桶的读写权限,同时桶区域需与Dataflow区域一致,避免跨区域访问问题。
    • 修改WriteToText的输出路径:Dataflow要求输出路径为目录而非具体文件名,建议改为gs://<bucket>/output/,WriteToText会自动生成带后缀的输出文件。
  3. 检查Dataflow服务账号权限

    • 确认Dataflow默认服务账号(格式为service-<project-number>@dataflow-service-producer-prod.iam.gserviceaccount.com)拥有以下权限:
      • Dataflow Worker角色
      • 目标GCS桶的Storage Object Creator/Editor权限
    • 确保项目已启用Dataflow API。
  4. 查看详细错误日志
    现有错误信息过于笼统,需登录GCP控制台进入对应Dataflow作业页面,查看Worker日志、作业详情中的错误事件,定位具体失败原因(如依赖缺失、资源不足、网络配置问题等)。

  5. 确认版本兼容性
    检查使用的Apache Beam版本是否与Dataflow服务兼容,建议使用官方推荐的稳定版本,避免版本不匹配导致的运行时异常。

内容的提问来源于stack exchange,提问作者luca

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 10:20:31