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

如何用Python高效读取Google Storage中的大量文件?

解决Google Cloud Storage批量读取数千文件的性能问题

问题根源

你遇到的Pickling client objects is not explicitly supported错误,是因为Google Cloud Storage的Blob对象无法被序列化(pickle),而多进程池Pool要求传递可序列化的参数,直接传递Blob对象会导致序列化失败。

可行解决方案

1. 传递Blob名称而非Blob对象(多进程改进版)

修改代码,先在主进程中获取所有Blob的名称字符串(字符串可序列化),子进程根据名称重新获取Blob对象。同时每个子进程独立创建Storage客户端,避免跨进程共享客户端引发的问题。

from google.cloud import storage
import time
import multiprocessing
from multiprocessing import Pool, Manager

cpu_count = multiprocessing.cpu_count()
manager = Manager()
finalized_list = manager.list()
bucket_name = "bucket-name"
cred_path = ".serviceAccountCredentials.json"

def list_blob_names():
    storage_client = storage.Client.from_service_account_json(cred_path)
    blobs = storage_client.list_blobs(bucket_name)
    return [blob.name for blob in blobs]

def read_blob_by_name(blob_name):
    client = storage.Client.from_service_account_json(cred_path)
    bucket = client.bucket(bucket_name)
    blob = bucket.blob(blob_name)
    with blob.open("r") as f:
        content = f.read()
        finalized_list.append(content)

def main():
    start_time = time.time()
    print(f"Start time: {start_time}")

    blob_names = list_blob_names()
    with Pool(processes=cpu_count) as pool:
        pool.map(read_blob_by_name, blob_names)
    
    end_time = time.time()
    elapsed_time = end_time - start_time
    print(f"Time taken: {elapsed_time} seconds")

if __name__ == "__main__":
    main()

2. 使用线程池(更适合IO密集型任务)

读取GCS文件属于IO密集型操作,线程池的开销远低于多进程——不需要跨进程复制数据,且GIL在IO等待时会自动释放,线程可并行处理任务。

from google.cloud import storage
import time
from concurrent.futures import ThreadPoolExecutor
import threading

bucket_name = "bucket-name"
cred_path = ".serviceAccountCredentials.json"
finalized_list = []
list_lock = threading.Lock()  # 线程锁保证列表操作线程安全

def list_blob_names():
    storage_client = storage.Client.from_service_account_json(cred_path)
    blobs = storage_client.list_blobs(bucket_name)
    return [blob.name for blob in blobs]

def read_blob_thread(blob_name):
    client = storage.Client.from_service_account_json(cred_path)
    bucket = client.bucket(bucket_name)
    blob = bucket.blob(blob_name)
    with blob.open("r") as f:
        content = f.read()
        with list_lock:
            finalized_list.append(content)

def main():
    start_time = time.time()
    print(f"Start time: {start_time}")

    blob_names = list_blob_names()
    # 线程数可设为CPU核心数的2-4倍,或根据文件数量调整
    with ThreadPoolExecutor(max_workers=cpu_count * 4) as executor:
        executor.map(read_blob_thread, blob_names)
    
    end_time = time.time()
    elapsed_time = end_time - start_time
    print(f"Time taken: {elapsed_time} seconds")

if __name__ == "__main__":
    main()

3. 额外优化建议

  • 批量预取元数据:获取Blob列表时批量拉取元数据,减少API调用次数。
  • 异步客户端:安装google-cloud-storage[aiohttp]使用异步客户端,结合asyncio实现更高并发效率,适合超大量文件场景。
  • 本地缓存:若后续需重复处理文件,可将内容缓存到本地磁盘或Redis,避免重复读取GCS。

内容的提问来源于stack exchange,提问作者Sajan Maharjan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 21:35:19