使用自定义模板创建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
相关产品推荐
相关产品推荐

