定义Pipeline Options后Dataflow作业停滞,求代码排查修正
问题
我编写了包含zip_extract、get_file_path、data_restructure三个方法的代码,预期执行流程为:
- 执行
zip_extract解压GCP Bucket中的zip文件; - 执行
get_file_path遍历Bucket获取所有文件路径列表; - 由
data_restructure根据文件是否为DICOM格式,将其存入目标Bucket的不同层级结构中。
我基于此编写的Dataflow管道代码如下:
with beam.Pipeline(options=pipeline_options) as p: file_paths = (p | "Get File Paths" >> beam.Create(get_file_path())) file_paths | "Data Restructure" >> beam.Map(lambda x: data_restructure(x))
但Dataflow日志报错:
The Dataflow job appears to be stuck because no worker activity has been seen in the last 1h. Please check the worker logs in Stackdriver Logging. You can also get help with Cloud Dataflow at https://cloud.google.com/dataflow/support.
完整代码:
def zip_extract(): ''' Function to unzip a folder in a bucket under a specific hierarchy ''' from google.cloud import storage client = storage.Client() bucket = client.bucket(landing_bucket) blobs_specific = list(bucket.list_blobs(prefix=data_folder)) for file_name in blobs_specific: file_extension = pathlib.Path(file_name.name).suffix try: if file_extension==".zip": destination_blob_pathname = file_name.name blob = 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) blob = bucket.blob(f'{file_name.name.replace(".zip","" )}/{contentfilename}') blob.upload_from_string(contentfile) logging.info("Unzip completed") except: logging.info('Skipping : {} file format found.'.format(file_extension)) continue client.close def get_file_path(): ''' Function to store all the file paths present in landing bucket into a list ''' zip_extract() file_paths = [] from google.cloud import storage client = storage.Client() bucket = client.bucket(landing_bucket) blobs_specific = list(bucket.list_blobs(prefix=data_folder)) try: for blob in blobs_specific: file_paths.append("gs://{}/".format(landing_bucket)+blob.name) client.close logging.info("List is ready with data") return file_paths except Exception as err: logging.error("Error while appending data to list : {}".format(err)) raise def data_restructure(line): ''' params line: String which has the file path Function to read each file and check if it is a DICOM file or not, if yes, store it in Study-Series-SOP hierarchy else store it in Descriptive folder in Intermediate bucket. ''' from google.cloud import storage InstanceUID={} client = storage.Client() destination_bucket = client.bucket(intermediate_bucket) cmd = "gsutil cp {} .\local_folder".format(line) result = subprocess.run(cmd,shell=True,capture_output=True,text=True) file_name=os.listdir(".\local_folder").pop(0) try: dicom_data = dcmread(".\local_folder\{}".format(file_name)) logging.info("Started reading Dicom file") for element in dicom_data: if element.name in ("Study Instance UID","Series Instance UID","SOP Instance UID","Modality"): InstanceUID[element.name]=element.value destination_bucket = client.bucket(intermediate_bucket) blob = destination_bucket.blob('Client/Test/DICOM/{}/{}/{}/{}.dcm'.format(list(InstanceUID.values())[1],list(InstanceUID.values())[2],list(InstanceUID.values())[3],list(InstanceUID.values())[0])) blob.upload_from_filename(".\local_folder\{}".format(file_name)) InstanceUID.clear() logging.info("DICOM file {} uploaded into Intermediate Bucket".format(file_name)) os.remove(".\local_folder\{}".format(file_name)) except Exception as e: file_extension = file_name.split("/")[-1].split(".")[-1] if file_extension != "zip" and "report" not in file_name and file_extension != "": blob = destination_bucket.blob('Test/Descriptive/{}'.format(file_name)) blob.upload_from_filename(".\local_folder\{}".format(file_name)) logging.info("Stored file into Descriptive folder") os.remove(".\local_folder\{}".format(file_name)) else: blob = destination_bucket.blob('Test/Reports/{}'.format(file_name)) blob.upload_from_filename(".\local_folder\{}".format(file_name)) logging.info("Stored Report file into Reports folder") os.remove(".\local_folder\{}".format(file_name)) client.close() def call_main(): parser = argparse.ArgumentParser() path_args, pipeline_args = parser.parse_known_args() pipeline_options = PipelineOptions(pipeline_args) setup_options= pipeline_options.view_as(SetupOptions) setup_options.setup_file='./setup.py' setup_options.save_main_session=True google_cloud_options = pipeline_options.view_as(GoogleCloudOptions) google_cloud_options.project = project_id google_cloud_options.job_name = "dataflow"+re.sub("[^0-9]+", "-", str(datetime.datetime.now())) google_cloud_options.service_account_email = "service_email" pipeline_options.view_as(StandardOptions).runner = "DataflowRunner" google_cloud_options.staging_location = config["staging_location"] google_cloud_options.temp_location = config["temp_location"] google_cloud_options.region = config["region"] pipeline_options.view_as(WorkerOptions).num_workers = 2 pipeline_options.view_as(WorkerOptions).machine_type = "n1-standard-2" pipeline_options.view_as(WorkerOptions).disk_size_gb = 1024 pipeline_options.view_as(WorkerOptions).network = vpc_name pipeline_options.view_as(WorkerOptions).subnetwork = f'regions/{config["region"]}/subnetworks/{subnet_name}' pipeline_options.view_as(WorkerOptions).use_public_ips=False with beam.Pipeline(options=pipeline_options) as p: file_paths = (p | "Get File Paths" >> beam.Create(get_file_path())) file_paths | "Data Restructure" >> beam.Map(lambda x: data_restructure(x)) if __name__ == '__main__': call_main()
setup.py文件:
import setuptools setuptools.setup( name='Installing Packages', version='1.0.0', install_requires=['google-cloud-datastore==1.15.3', 'google.cloud.storage==1.16.1', 'apache-beam[gcp]==2.31.0', 'google-api-core==1.33.2', 'google-cloud-core==1.7.3', 'google-cloud-logging == 1.15.1', 'pydicom == 2.3.1', 'uuid == 1.30', 'google-cloud-secret-manager', 'psycopg2-binary'], packages=setuptools.find_packages())
我是Apache Beam和Dataflow新手,尝试过多种管道写法均未解决问题,恳请指出代码中的错误,说明当前转换写法是否正确,并提供正确实现方案。
问题分析
当前代码存在多个核心错误,导致Dataflow作业卡死:
beam.Create(get_file_path())执行位置错误:get_file_path()会在提交作业的本地机器上执行,而非Dataflow Worker。zip_extract()被其调用也会在本地跑,导致Worker上没有待处理的文件路径,管道无数据流动,最终卡住。data_restructure本地操作风险:- 依赖
gsutil命令,但Worker环境不一定预装,且.\local_folder可能不存在; - 用
os.listdir().pop(0)获取文件不可靠,多Worker并发时会出现文件冲突或错误。
- 依赖
- 资源未正确释放:
client.close()只是引用方法,未执行(缺少括号client.close()),会导致连接泄漏。 - Pipeline转换逻辑不符合分布式模型:用本地生成的列表作为数据源,违背Dataflow分布式执行的设计,无法利用Worker集群能力。
正确实现方案
1. 重构Pipeline核心逻辑
将所有操作放到Dataflow分布式转换中,确保逻辑在Worker上执行,并通过信号控制流程顺序:
def call_main(): # 保留原有的pipeline_options配置(略) with beam.Pipeline(options=pipeline_options) as p: # 步骤1:获取Bucket中所有Blob路径 blobs = (p | "Init Blob List" >> beam.Create([data_folder]) | "Fetch Blob Paths" >> beam.FlatMap(list_bucket_blobs)) # 步骤2:仅处理ZIP文件,完成解压 unzip_results = (blobs | "Filter ZIPs" >> beam.Filter(lambda x: x.endswith(".zip")) | "Unzip Blobs" >> beam.Map(unzip_blob)) # 步骤3:等待解压完成后,重新获取所有非ZIP文件路径 all_files = (p | "Trigger Post-Unzip" >> beam.Create([data_folder]) | "Wait For Unzip" >> beam.WaitOnSignal(unzip_results) | "Fetch All Files" >> beam.FlatMap(list_bucket_blobs) | "Exclude ZIPs" >> beam.Filter(lambda x: not x.endswith(".zip"))) # 步骤4:分布式处理每个文件 all_files | "Process Files" >> beam.Map(process_single_file)
2. 工具函数重构
列表Bucket文件的函数
def list_bucket_blobs(prefix): from google.cloud import storage client = storage.Client() bucket = client.bucket(landing_bucket) blob_paths = [] for blob in bucket.list_blobs(prefix=prefix): blob_paths.append(f"gs://{landing_bucket}/{blob.name}") client.close() return blob_paths
解压Blob的函数
def unzip_blob(zip_blob_path): from google.cloud import storage from zipfile import ZipFile, is_zipfile import io import pathlib client = storage.Client() bucket = client.bucket(landing_bucket) zip_blob_name = zip_blob_path.replace(f"gs://{landing_bucket}/", "") blob = bucket.blob(zip_blob_name) zipbytes = io.BytesIO(blob.download_as_bytes()) if is_zipfile(zipbytes): with ZipFile(zipbytes, 'r') as myzip: for contentfilename in myzip.namelist(): contentfile = myzip.read(contentfilename) dest_blob_name = f'{zip_blob_name.replace(".zip", "")}/{contentfilename}' dest_blob = bucket.blob(dest_blob_name) dest_blob.upload_from_string(contentfile) client.close() return f"Unzipped: {zip_blob_path}"
处理单个文件的函数
def process_single_file(file_path): from google.cloud import storage from pydicom import dcmread import tempfile import os client = storage.Client() src_bucket = client.bucket(landing_bucket) dest_bucket = client.bucket(intermediate_bucket) # 使用临时文件避免本地目录冲突 with tempfile.NamedTemporaryFile(delete=False) as tmp_file: src_blob_name = file_path.replace(f"gs://{landing_bucket}/", "") src_blob = src_bucket.blob(src_blob_name) src_blob.download_to_file(tmp_file) tmp_file_path = tmp_file.name try: # 尝试读取DICOM文件 dicom_data = dcmread(tmp_file_path) instance_uid = { "Study": dicom_data.get("StudyInstanceUID", ""), "Series": dicom_data.get("SeriesInstanceUID", ""), "SOP": dicom_data.get("SOPInstanceUID", ""), "Modality": dicom_data.get("Modality", "") } # 构造DICOM目标路径 dest_blob_name = f'Client/Test/DICOM/{instance_uid["Study"]}/{instance_uid["Series"]}/{instance_uid["SOP"]}/{instance_uid["Modality"]}.dcm' dest_blob = dest_bucket.blob(dest_blob_name) dest_blob.upload_from_filename(tmp_file_path) except Exception as e: # 处理非DICOM文件 file_name = os.path.basename(file_path) file_ext = os.path.splitext(file_name)[1].lower() if file_ext != ".zip" and "report" not in file_name.lower() and file_ext: dest_blob_name = f'Test/Descriptive/{file_name}' else: dest_blob_name = f'Test/Reports/{file_name}' dest_blob = dest_bucket.blob(dest_blob_name) dest_blob.upload_from_filename(tmp_file_path) # 清理临时文件 os.unlink(tmp_file_path) client.close()
3. setup.py 依赖调整(保证版本兼容)
import setuptools setuptools.setup( name='dicom-dataflow-processing', version='1.0.0', install_requires=[ 'apache-beam[gcp]==2.31.0', 'google-cloud-storage==1.44.0', # 与beam 2.31兼容的版本 'pydicom==2.3.1', 'google-cloud-secret-manager==2.10.0', 'psycopg2-binary==2.9.5' ], packages=setuptools.find_packages() )
关键修正点说明
- 分布式执行:所有Bucket操作、解压、文件处理都在Dataflow Worker上执行,避免本地执行导致的Worker无数据问题。
- 临时文件管理:用
tempfile替代固定本地目录,解决多Worker并发时的文件冲突。 - 依赖兼容性:调整
google-cloud-storage版本,确保与Apache Beam 2.31.0兼容。 - 资源释放:正确调用
client.close()释放连接。 - 顺序控制:用
beam.WaitOnSignal保证解压完成后再获取文件列表,避免遗漏解压后的文件。
内容的提问来源于stack exchange,提问作者user13109553
相关产品推荐
相关产品推荐

