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

使用自定义模板创建Dataflow任务时遇无步骤指定错误求助

Apache Beam Dataflow报错:"Runnable workflow has no steps specified"排查

我尝试用自定义模板创建Apache Beam Dataflow任务时,收到错误提示:"Runnable workflow has no steps specified",日志里没有其他相关信息。已经创建了虚拟环境,执行的代码如下:

import sys
import os
import apache_beam as beam
import google.cloud.logging
import google.auth
from google.cloud import storage
import pandas as pd
import io
import gcsfs as gcs
from io import BytesIO
import re
from datetime import date
import datetime
import logging
import argparse
from pyspark.context import SparkContext
from pyspark import SparkConf
from pyspark import SparkConf, SparkContext
from pyspark.sql import SparkSession
from pyspark.sql.types import *
from pyspark import SparkConf, SparkContext
from py4j.java_gateway import java_import
from pyspark.sql.functions import udf
from apache_beam.io import ReadFromText
from apache_beam.io import WriteToText
from apache_beam.options.pipeline_options import PipelineOptions
from apache_beam.options.pipeline_options import SetupOptions
import urllib
import urllib.request
from urllib.request import urlretrieve
from zipfile import ZipFile, is_zipfile


date = datetime.datetime.now()
monthname = 'September'
yearname = date.strftime("%Y")

print(monthname)
print(yearname)

url = 'https://www.abc.gov/files/zip/statecontract-'+monthname+'-'+yearname+'-employee.zip'

destination_zip_name = 'upload.zip'

def upload_blob(bucket_name, url, destination_zip_name,argv=None, save_main_session=True):

    parser = argparse.ArgumentParser()
    parser.add_argument(
        '--input',
        dest='input',
        help='Input file to process.')
    parser.add_argument(
        '--output',
        dest='output',
        help='Output file to write results to.')
    known_args, pipeline_args = parser.parse_known_args(argv)

    pipeline_options = PipelineOptions(pipeline_args)
    pipeline_options.view_as(SetupOptions).save_main_session = save_main_session

    with beam.Pipeline(options=pipeline_options) as p:
        storage_client = storage.Client()
        source_bucket = storage_client.get_bucket(bucket_name)
        print('source bucket - ',source_bucket)
        destination_bucket_name = storage_client.get_bucket(bucket_name)
        print('destination bucket - ',destination_bucket_name)


        my_file = urllib.request.urlopen(url)
        blob1 = source_bucket.blob(destination_zip_name)
        blob1.upload_from_string(my_file.read(), content_type='application/zip')

        destination_blob_pathname = destination_zip_name
        print('destination_blob_pathname - ',destination_blob_pathname)

        blob = source_bucket.blob(destination_blob_pathname)
        zipbytes = io.BytesIO(blob.download_as_string())
        

        if is_zipfile(zipbytes):
            with ZipFile(zipbytes, 'r') as myzip:
                for contentfilename in myzip.namelist():
                    contentfile = myzip.read(contentfilename)
                    
                    if '.csv' in contentfilename.casefold():
                        output_file = f'/tmp/{contentfilename.split("/")[-1]}'
                        print('output_file - ',output_file)
                        outfile = open(output_file, 'wb')
                        outfile.write(contentfile)
                        outfile.close()

                        blob = source_bucket.blob(
                            f'{destination_zip_name.rstrip(".zip")}/{contentfilename}'
                        )
                        with open(output_file, "rb") as my_csv:
                            blob.upload_from_file(my_csv)

        blob1.delete()                
        print('done running function')


if __name__ == '__main__':
    upload_blob('testbucket', url, destination_zip_name)

问题根源

你的代码虽然初始化了beam.Pipeline上下文,但没有向Pipeline中添加任何Beam标准的数据流转换(Transform)步骤。Dataflow要求工作流必须包含至少一个可执行的Beam操作(比如beam.Create、beam.ParDo等),否则会判定工作流无有效步骤,抛出该错误。

当前所有业务逻辑(下载文件、GCS上传、解压)都是在beam.Pipeline()上下文内执行的普通Python代码,没有使用Beam的PCollection和Transform体系,相当于Pipeline是空的。


解决方案

根据需求选择以下两种方案之一:

方案1:改为普通Python脚本(无需Dataflow)

如果不需要分布式处理能力,直接移除Beam相关代码,用普通脚本执行逻辑:

import google.cloud.storage
import urllib.request
import io
from zipfile import ZipFile, is_zipfile
import datetime

date = datetime.datetime.now()
monthname = 'September'
yearname = date.strftime("%Y")

url = f'https://www.abc.gov/files/zip/statecontract-{monthname}-{yearname}-employee.zip'
destination_zip_name = 'upload.zip'

def upload_blob(bucket_name, url, destination_zip_name):
    storage_client = google.cloud.storage.Client()
    source_bucket = storage_client.get_bucket(bucket_name)
    
    # 下载并上传zip文件到GCS
    my_file = urllib.request.urlopen(url)
    blob1 = source_bucket.blob(destination_zip_name)
    blob1.upload_from_string(my_file.read(), content_type='application/zip')
    
    # 下载zip文件并解压上传CSV
    zipbytes = io.BytesIO(blob1.download_as_string())
    if is_zipfile(zipbytes):
        with ZipFile(zipbytes, 'r') as myzip:
            for contentfilename in myzip.namelist():
                if '.csv' in contentfilename.casefold():
                    contentfile = myzip.read(contentfilename)
                    blob = source_bucket.blob(f'{destination_zip_name.rstrip(".zip")}/{contentfilename}')
                    blob.upload_from_string(contentfile)
    
    # 删除临时zip文件
    blob1.delete()                
    print('done running function')

if __name__ == '__main__':
    upload_blob('testbucket', url, destination_zip_name)

方案2:重构为标准Beam Pipeline(适配Dataflow)

如果需要利用Dataflow的分布式能力,将逻辑封装为Beam的DoFn,构建完整的数据流:

import apache_beam as beam
import google.cloud.storage
import urllib.request
import io
from zipfile import ZipFile, is_zipfile
import datetime
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions
import argparse

date = datetime.datetime.now()
monthname = 'September'
yearname = date.strftime("%Y")

url = f'https://www.abc.gov/files/zip/statecontract-{monthname}-{yearname}-employee.zip'
destination_zip_name = 'upload.zip'

class ProcessZipFile(beam.DoFn):
    def process(self, element):
        bucket_name, url, dest_zip = element
        storage_client = google.cloud.storage.Client()
        source_bucket = storage_client.get_bucket(bucket_name)
        
        # 下载并上传zip到GCS
        my_file = urllib.request.urlopen(url)
        blob1 = source_bucket.blob(dest_zip)
        blob1.upload_from_string(my_file.read(), content_type='application/zip')
        
        # 解压并上传CSV文件
        zipbytes = io.BytesIO(blob1.download_as_string())
        if is_zipfile(zipbytes):
            with ZipFile(zipbytes, 'r') as myzip:
                for contentfilename in myzip.namelist():
                    if '.csv' in contentfilename.casefold():
                        contentfile = myzip.read(contentfilename)
                        blob = source_bucket.blob(f'{dest_zip.rstrip(".zip")}/{contentfilename}')
                        blob.upload_from_string(contentfile)
        
        blob1.delete()
        yield "Processing completed"

def run(argv=None, save_main_session=True):
    parser = argparse.ArgumentParser()
    parser.add_argument('--bucket', dest='bucket', required=True, help='GCS bucket name')
    known_args, pipeline_args = parser.parse_known_args(argv)
    
    pipeline_options = PipelineOptions(pipeline_args)
    pipeline_options.view_as(SetupOptions).save_main_session = save_main_session
    
    with beam.Pipeline(options=pipeline_options) as p:
        (p 
         | 'Initialize Task' >> beam.Create([(known_args.bucket, url, destination_zip_name)])
         | 'Process Zip File' >> beam.ParDo(ProcessZipFile())
         | 'Log Result' >> beam.Map(print)
        )

if __name__ == '__main__':
    run(argv=['--bucket', 'testbucket'])

关键说明

  • Beam Pipeline的核心是PCollection(数据集) + Transform(转换操作),必须包含至少一组这样的组合,Dataflow才会识别为有效工作流。
  • 简单文件操作优先选择方案1;需要分布式处理大量文件时,方案2更适合。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 15:10:23