使用Pandas read_csv多进程读取海量CSV的内存问题及优化咨询
问题描述
在Sagemaker的ml.m5.4xlarge实例(16vCPU、64GB内存)上,使用以下代码读取并合并数千个CSV文件为大型DataFrame,供XGBoost训练:
def _read_training_data(training_data_path: str) -> pd.DataFrame: df = pd.read_csv(training_data_path) return df def read_training_data( paths: List[str] ) -> pd.DataFrame: # lazy loading of modules for training import multiprocessing as mp from concurrent.futures import ProcessPoolExecutor, as_completed # get directories ipaths = train_mdirnames(paths) logger.info(f'start parallel data reading with {mp.cpu_count()} core') df = None with ProcessPoolExecutor(max_workers=mp.cpu_count()) as executor: tasks = [ executor.submit(_read_training_data, ipath) for ipath in ipaths] for future in as_completed(tasks): try: _df = future.result() df = _df if df is None else pd.concat([df, _df]) except Exception as e: raise e logger.info(f'have read {len(df)} data points')
处理单日400万行(6GB)数据时耗时约11小时并成功完成,但处理一周量级(约41GB)数据时出现内存溢出问题。需要找到海量数据读取的最优方案,以在同类型或更小实例上处理更大的月度数据,同时将合并后的DataFrame送入XGBoost训练。
解决方案
针对内存溢出和效率问题,可从以下几个方向优化:
1. 优化单文件读取的内存占用
Pandas默认读取CSV时自动推断数据类型,易占用冗余内存。在_read_training_data中添加参数限制内存:
def _read_training_data(training_data_path: str) -> pd.DataFrame: # 提前定义各列最优数据类型,用小精度类型替代大类型 dtype_map = { "count_col": "int16", "score_col": "float32", "type_col": "category" } # 只读取训练必需的列,跳过冗余数据 usecols = ["count_col", "score_col", "type_col", "label"] df = pd.read_csv( training_data_path, dtype=dtype_map, usecols=usecols, low_memory=False # 避免类型推断时的内存波动 ) return df
若单文件过大,可通过chunksize分块读取合并:
def _read_training_data(training_data_path: str) -> pd.DataFrame: dtype_map = {...} usecols = [...] chunk_list = [] # 按10万行分块读取,减少单块内存占用 for chunk in pd.read_csv( training_data_path, dtype=dtype_map, usecols=usecols, chunksize=100000 ): chunk_list.append(chunk) return pd.concat(chunk_list, ignore_index=True)
2. 改进并行合并逻辑
当前代码每次pd.concat([df, _df])都会生成新DataFrame,导致内存中同时存在新旧两个大对象,加剧内存消耗。改为先收集所有小DataFrame到列表,最后一次性合并:
def read_training_data(paths: List[str]) -> pd.DataFrame: import multiprocessing as mp from concurrent.futures import ProcessPoolExecutor, as_completed ipaths = train_mdirnames(paths) logger.info(f'start parallel data reading with {mp.cpu_count()} cores') df_list = [] with ProcessPoolExecutor(max_workers=mp.cpu_count()) as executor: tasks = {executor.submit(_read_training_data, ipath): ipath for ipath in ipaths} for future in as_completed(tasks): try: _df = future.result() df_list.append(_df) except Exception as e: logger.error(f'Failed to read {tasks[future]}: {e}') raise e # 一次性合并所有小DataFrame,避免多次复制内存 df = pd.concat(df_list, ignore_index=True) logger.info(f'have read {len(df)} data points') return df
3. 用Dask替代Pandas处理超大数据
Dask将数据拆分为多个分区并行处理,无需全量加载到内存,完美适配海量数据场景,且可直接对接XGBoost训练:
import dask.dataframe as dd from dask.distributed import Client def read_training_data_dask(paths: List[str]) -> dd.DataFrame: # 初始化Dask客户端,利用实例全部CPU核心 client = Client(n_workers=mp.cpu_count()) # 读取所有CSV为Dask DataFrame,自动分区 dtype_map = {...} usecols = [...] ddf = dd.read_csv( paths, dtype=dtype_map, usecols=usecols, blocksize="64MB" # 按64MB分区,可根据内存调整 ) return ddf
训练时直接传入Dask DataFrame:
import xgboost as xgb # 拆分特征和标签 X = ddf.drop("label", axis=1) y = ddf["label"] # 使用Dask版XGBoost训练 model = xgb.dask.DaskXGBClassifier(n_estimators=100) model.fit(X, y)
4. 转换数据格式为Parquet
CSV是文本格式,读取慢且内存占用高,转换为Parquet列式存储可大幅优化:
- 存储空间仅为CSV的1/3~1/5
- 读取速度提升数倍
- 自动保留数据类型,无需重复推断
转换与读取示例:
# 批量转换CSV为Parquet def convert_csv_to_parquet(csv_paths: List[str], parquet_dir: str): import os os.makedirs(parquet_dir, exist_ok=True) for csv_path in csv_paths: df = _read_training_data(csv_path) parquet_path = os.path.join(parquet_dir, os.path.basename(csv_path).replace(".csv", ".parquet")) df.to_parquet(parquet_path, engine="pyarrow") # 读取Parquet文件 def _read_parquet_data(parquet_path: str) -> pd.DataFrame: return pd.read_parquet(parquet_path, engine="pyarrow")
5. Sagemaker管道训练模式
若无需全量合并数据再训练,可使用Sagemaker管道模式(Pipe Mode),边读取数据边训练,避免全量加载的内存压力。XGBoost的Sagemaker容器支持直接从S3读取管道数据,无需将所有数据加载到实例内存。
内容的提问来源于stack exchange,提问作者ethicalguy
相关产品推荐
相关产品推荐

