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

Python多进程拆分40G大CSV按列存储无效,求排查方案

问题根源

你的代码里每个进程都完整读取了40G的CSV文件,这才是导致看起来像串行的核心原因。磁盘IO是绝对瓶颈,多个进程同时抢读同一个大文件,会把磁盘负载拉满,每个进程的读取速度都会变慢,最终整体效率还不如串行,甚至可能因为内存不足拖垮系统。

多进程写入文件本身是可行的,你的问题完全出在错误的执行逻辑上——把最耗时的读文件环节重复执行了N次。

正确处理方案

1. 内存足够时的一次性读取方案

先在主进程把整个文件读进来,再拆分列分配给多进程处理,避免重复读大文件:

import pandas as pd
import os
from multiprocessing import Pool

def process_column(col_data, output_path, col_name):
    # 处理单列:删除全空行
    cleaned_data = col_data.dropna(how='all')
    # 写入单独CSV
    cleaned_data.to_csv(os.path.join(output_path, f'{col_name}.csv'), index=False)

if __name__ == '__main__':
    original_large_data_path = '你的大文件路径.csv'
    output_path = '输出目录路径'
    columns_to_split = ['列名1', '列名2', ...]  # 替换成你要拆分的目标列

    # 主进程一次性读取目标列(减少内存占用)
    df = pd.read_csv(original_large_data_path, usecols=columns_to_split)

    # 确保输出目录存在
    os.makedirs(output_path, exist_ok=True)

    # 准备进程任务参数
    tasks = [(df[col], output_path, col) for col in columns_to_split]

    # 启动进程池处理
    with Pool(5) as pool:
        pool.starmap(process_column, tasks)

2. 内存不足时的分块读取方案

如果40G文件直接读入内存扛不住,就用分块读取,把每一块的对应列数据追加到输出文件:

import pandas as pd
import os
from multiprocessing import Pool
from functools import partial

def process_chunk(chunk, output_path, col_name):
    col_data = chunk[col_name].dropna(how='all')
    # 追加模式写入,第一次写加表头,后续不加
    file_path = os.path.join(output_path, f'{col_name}.csv')
    mode = 'a' if os.path.exists(file_path) else 'w'
    col_data.to_csv(file_path, mode=mode, header=(mode == 'w'), index=False)

if __name__ == '__main__':
    original_large_data_path = '你的大文件路径.csv'
    output_path = '输出目录路径'
    columns_to_split = ['列名1', '列名2', ...]
    chunksize = 100000  # 每次读10万行,根据内存调整

    os.makedirs(output_path, exist_ok=True)

    with Pool(5) as pool:
        # 给每个列绑定处理参数,生成偏函数
        for col in columns_to_split:
            process_func = partial(process_chunk, output_path=output_path, col_name=col)
            # 分块读取并提交任务
            for chunk in pd.read_csv(original_large_data_path, usecols=[col], chunksize=chunksize):
                pool.apply_async(process_func, args=(chunk,))
        pool.close()
        pool.join()
额外注意事项
  • 进程数别设得比CPU核心数太多,磁盘IO是瓶颈,太多进程只会互相抢资源。
  • 如果CSV有特殊格式(比如引号、自定义分隔符),读文件时要加上sep、quotechar等参数避免解析错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 04:27:27