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

Python并行处理结果如何汇总至单个CSV?解决全局变量失效问题

高效并行处理CSV并汇总结果的方案

核心逻辑

别用全局变量,让每个并行任务直接返回处理后的结果数据,主进程收集所有返回值后统一合并成DataFrame再写入CSV——全程只做两次IO操作,比生成一堆小文件再合并高效太多,还能避开多进程内存隔离导致的全局变量失效问题。


方案1:用Python标准库concurrent.futures

无需额外安装,CPU密集型任务首选:

import pandas as pd
from concurrent.futures import ProcessPoolExecutor

def process_row(row):
    # 替换成你的实际处理逻辑,比如字段计算、数据转换
    return {
        'id': row['id'],
        'processed_value': row['original_value'] * 2  # 示例操作
    }

if __name__ == '__main__':
    # 读取输入CSV
    input_df = pd.read_csv('input.csv')
    # 转成字典列表,方便逐行传递给并行任务
    rows = input_df.to_dict('records')
    
    # 启动进程池处理
    with ProcessPoolExecutor() as executor:
        # 用map分发任务,自动收集所有返回结果
        results = list(executor.map(process_row, rows))
    
    # 合并结果并写入最终CSV
    output_df = pd.DataFrame(results)
    output_df.to_csv('output.csv', index=False)

方案2:用joblib(如果你习惯这个库)

和上面逻辑一致,只是调用方式不同:

import pandas as pd
from joblib import Parallel, delayed

def process_row(row):
    # 替换为你的处理逻辑
    return {
        'id': row['id'],
        'processed_value': row['original_value'] * 2
    }

if __name__ == '__main__':
    input_df = pd.read_csv('input.csv')
    rows = input_df.to_dict('records')
    
    # n_jobs=-1表示用满所有CPU核心
    results = Parallel(n_jobs=-1)(delayed(process_row)(row) for row in rows)
    
    output_df = pd.DataFrame(results)
    output_df.to_csv('output.csv', index=False)

方案3:超大数据量优化(分块处理)

如果CSV大到内存装不下,就分块读+并行处理:

import pandas as pd
from concurrent.futures import ProcessPoolExecutor

def process_chunk(chunk):
    # 直接处理整个数据块,返回处理后的DataFrame
    chunk['processed_value'] = chunk['original_value'] * 2  # 示例操作
    # 只保留需要的字段,减少内存占用
    return chunk[['id', 'processed_value']]

if __name__ == '__main__':
    # chunksize根据你的内存情况调整,比如1000行一块
    chunk_iter = pd.read_csv('input.csv', chunksize=1000)
    
    with ProcessPoolExecutor() as executor:
        processed_chunks = list(executor.map(process_chunk, chunk_iter))
    
    # 合并所有处理后的块
    output_df = pd.concat(processed_chunks, ignore_index=True)
    output_df.to_csv('output.csv', index=False)

为什么你的全局变量方案没用?

Python多进程里,每个子进程会复制主进程的内存空间,子进程改全局变量只改自己副本里的,主进程根本看不到,所以最后output_df是空的。上面的方案都是让子进程把结果return出来,主进程主动收集,从根上解决了这个问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 04:55:32