Python中使用Pub/Sub订阅时GCS文件下载缓慢问题求助
问题分析与解决方案
核心问题原因
- 回调函数错误触发线程循环:
callback函数末尾调用了task_manager(),导致每收到一条Pub/Sub消息就启动一个新的无限循环线程。随着消息增多,线程数量暴增,引发严重的CPU上下文切换开销,直接拖慢所有下载任务。 - GCS客户端重复初始化:
download_folder函数每次都新建storage.Client(),重复执行认证和连接建立流程,额外增加了延迟。 - Pub/Sub订阅线程被挤占:默认情况下Pub/Sub订阅者的线程池资源有限,错误的线程逻辑会占用订阅线程,导致消息处理与下载任务互相抢占资源。
修复步骤
1. 移除回调中的task_manager()调用
脚本开头已经单独启动了task_manager线程,无需在每次回调中重复调用,否则会创建大量重复的无限循环线程。
2. 复用GCS客户端
将GCS客户端初始化移到全局作用域,避免每次下载都重复创建和认证,减少连接开销。
3. 优化线程模型(可选)
使用concurrent.futures.ThreadPoolExecutor处理下载任务,替代单线程的task_manager,提升多任务并发效率;同时显式配置Pub/Sub订阅者的线程池大小,避免订阅线程被阻塞。
修改后的完整代码
from google.cloud import pubsub_v1, storage import os import queue import threading from concurrent.futures import ThreadPoolExecutor # Setting up authentication os.environ['GOOGLE_APPLICATION_CREDENTIALS'] = 'key-path' # 复用GCS客户端,避免重复初始化 gcs_client = storage.Client() # Initializing Pub/Sub subscriber subscriber = pubsub_v1.SubscriberClient() subscription_path = subscriber.subscription_path('project-name', 'subscription-name') # Setting up a task queue task_queue = queue.Queue() # 配置线程池处理下载任务,可根据服务器资源调整线程数 download_executor = ThreadPoolExecutor(max_workers=4) def download_folder(bucket_name, source_folder_name, destination_folder_path): bucket = gcs_client.get_bucket(bucket_name) blobs = bucket.list_blobs(prefix=f"{source_folder_name}/") print("Authorization DONE") for blob in blobs: if blob.name.endswith('/'): continue # Skipping directories print("FIRST FILE DOWNLOADING") # Ensuring the destination directory exists os.makedirs(destination_folder_path, exist_ok=True) destination_file = f"{destination_folder_path}/{blob.name.split('/')[-1]}" print(destination_file) # Downloading the blob to a file blob.download_to_filename(destination_file) print(f"Downloaded {blob.name}") def process_task(message_list): print("Start downloading media files....") source_folder_name = message_list[0] destination_folder_path = f'D:/{source_folder_name}' print(source_folder_name) print(destination_folder_path) # 用线程池执行下载任务,提升并发效率 download_executor.submit(download_folder, "bucket-name", source_folder_name, destination_folder_path) print(f'Processed task with data: {message_list}') def task_manager(): while True: next_task_data = task_queue.get() process_task(next_task_data) task_queue.task_done() # 标记任务完成,避免队列积压 def callback(message): print(f'Received new message') message_string = message.data.decode('utf-8') message_list = message_string.split(',') print(message_string) # Adding the message data to the task queue task_queue.put(message_list) message.ack() # 移除错误的task_manager()调用 # Starting the task manager in a separate thread processor_thread = threading.Thread(target=task_manager, daemon=True) processor_thread.start() # 配置Pub/Sub订阅者的线程池,避免线程不足 subscriber_options = pubsub_v1.types.SubscriberOptions( flow_control=pubsub_v1.types.FlowControl(max_messages=10), scheduler=pubsub_v1.types.SchedulerOptions(worker_threads=4) ) # Subscribing to the Pub/Sub topic future = subscriber.subscribe(subscription_path, callback=callback, subscriber_options=subscriber_options) try: print('Listening for messages...') future.result() # 用future.result()替代无限循环,更优雅 except KeyboardInterrupt: future.cancel() download_executor.shutdown(wait=True)
额外优化建议
- 增加错误处理:在下载任务中添加try-except块,避免单个下载失败导致整个task_manager线程崩溃。
- 配置Flow Control:根据自身服务器资源调整Pub/Sub的Flow Control参数,避免消息堆积过多。
- 日志替代print:使用Python的
logging模块替代print,便于生产环境排查问题。
内容的提问来源于stack exchange,提问作者vlad Hohlov
相关产品推荐
相关产品推荐

