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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 20:44:51