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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 13:14:53