GCP云存储多CSV文件快速读取合并及convtools报错解决
GCP云存储多CSV文件高效读取方案(解决pandas慢、convtools文件不存在问题)
一、解决convtools的"文件不存在"错误
convtools的Table.from_csv默认无法直接识别gs://协议路径,会将其当作本地文件查找,因此抛出文件不存在错误。需通过GCS客户端获取文件流后传入:
修改后的convtools代码:
from convtools import conversion as c from convtools.contrib.tables import Table from google.cloud import storage client = storage.Client() bucket = client.get_bucket(bucket_name) file_list = list_blobs_with_prefix( bucket_name=bucket_name, prefix=prefix + date_timestamp ) table = None for blob_name in file_list: blob = bucket.blob(blob_name) # 打开GCS文件流 with blob.open("r") as f: table_ = Table.from_csv(f, header=True) if table is None: table = table_ else: table.chain(table_) iterable_dict = table.into_iter_rows(dict) df = pd.DataFrame(list(iterable_dict))
二、更高效的多CSV读取优化方案
1. gcsfs + Dask并行读取(推荐)
Dask支持并行加载多文件,配合gcsfs可直接解析gs://路径,大幅提升大文件集的读取速度:
import dask.dataframe as dd import gcsfs # 初始化GCS文件系统 fs = gcsfs.GCSFileSystem() file_list = list_blobs_with_prefix( bucket_name=bucket_name, prefix=prefix + date_timestamp ) file_paths = [f"gs://{bucket_name}/{blob}" for blob in file_list] # 并行读取所有CSV文件 ddf = dd.read_csv(file_paths) # 转换为pandas DataFrame(数据量过大时建议保留Dask格式处理) df = ddf.compute()
2. 多进程优化pandas原生读取
在pandas生态内,通过多进程利用CPU多核,提升单进程concat的效率:
import pandas as pd from google.cloud import storage from multiprocessing import Pool client = storage.Client() bucket = client.get_bucket(bucket_name) def read_blob_to_df(blob_name): blob = bucket.blob(blob_name) with blob.open("r") as f: return pd.read_csv(f) file_list = list_blobs_with_prefix( bucket_name=bucket_name, prefix=prefix + date_timestamp ) # 多进程并行读取文件 with Pool() as pool: dfs = pool.map(read_blob_to_df, file_list) df = pd.concat(dfs, ignore_index=True)
3. 临时目录批量下载读取(适合小文件集)
利用Cloud Functions临时存储空间批量下载文件后读取,注意临时空间容量限制:
import pandas as pd from google.cloud import storage import tempfile import os client = storage.Client() bucket = client.get_bucket(bucket_name) file_list = list_blobs_with_prefix( bucket_name=bucket_name, prefix=prefix + date_timestamp ) dfs = [] with tempfile.TemporaryDirectory() as tmpdir: for blob_name in file_list: local_path = os.path.join(tmpdir, os.path.basename(blob_name)) bucket.blob(blob_name).download_to_filename(local_path) dfs.append(pd.read_csv(local_path)) df = pd.concat(dfs, ignore_index=True)
内容的提问来源于stack exchange,提问作者erik25
相关产品推荐
相关产品推荐

