如何用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
相关产品推荐
相关产品推荐

