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

如何高效将Python ProcessPoolExecutor结果写入单个Parquet文件?

解决大型CSV转Parquet时的内存峰值问题:边处理边写入

问题根源

你的当前实现中,executor.map会等待所有进程完成转换后才返回结果集合,随后pd.concat会将所有转换后的DataFrame合并成一个超大对象,这直接导致内存占用飙升。

解决方案:边处理边写入

核心思路是一有转换完成的chunk就立即写入Parquet文件,无需等待所有任务结束,也不需要合并成大DataFrame。关键是用concurrent.futures.as_completed跟踪完成的任务,在主进程中逐个处理写入(ParquetWriter不是线程/进程安全的,必须由主进程操作)。

修改后的完整代码

import pyarrow as pa
import pyarrow.parquet as pq
import pandas as pd
import regex as re
from concurrent.futures import ProcessPoolExecutor, as_completed
from multiprocessing import cpu_count

filename = 'input.csv'

# 生成示例输入文件
n_rows = int(25 * 1e6)
input_df = pd.DataFrame([f'User_id_{i}_says_"hello!".' for i in range(n_rows)], columns=['A'])
input_df.to_csv(filename, header=True, index=False)

# 文本清理函数
def clean_text(s):
    regex = re.compile('[^a-zA-Z]')
    return regex.sub('', s)

# 单chunk转换逻辑
def process(df):
    out = pd.DataFrame(df['A'].str.split('_').to_list())
    out[2] = out[2].astype(str)
    out[4] = out[4].apply(clean_text)
    out.rename(columns={i: f'A{i}' for i in out.columns}, inplace=True)
    return out

if __name__ == "__main__":
    # 定义Parquet schema
    schema = pa.schema([(f'A{i}', pa.string()) for i in range(5)])
    
    core_count = cpu_count() - 3
    chunk_size = int(1e5)
    
    # 读取CSV分块
    chunks = pd.read_csv('input.csv', header=0, na_filter=False, chunksize=chunk_size)
    
    # 初始化Parquet写入器(主进程操作)
    with pq.ParquetWriter('processed.parquet', schema=schema, compression='GZIP') as writer:
        # 提交所有转换任务到进程池
        with ProcessPoolExecutor(core_count) as executor:
            # 用字典保存future和对应的chunk(可选,用于异常追踪)
            future_to_chunk = {executor.submit(process, chunk): chunk for chunk in chunks}
            
            # 遍历完成的任务,立即写入
            for future in as_completed(future_to_chunk):
                try:
                    transformed_df = future.result()
                    # 转换为Arrow RecordBatch并写入
                    batch = pa.RecordBatch.from_pandas(transformed_df, schema=schema)
                    writer.write_batch(batch)
                except Exception as e:
                    print(f"处理chunk时出错: {e}")

关键改动说明

  1. 用as_completed替代map:不再等待所有任务完成,而是实时获取已完成的转换结果
  2. 主进程负责写入:ParquetWriter由主进程创建和调用,避免多进程共享资源的安全问题
  3. 无合并操作:每个转换后的chunk直接写入,内存仅保留当前chunk的数据,内存占用稳定

额外优化建议

  • 可以考虑在process函数中直接返回Arrow RecordBatch,减少主进程的转换开销:
    def process(df):
        out = pd.DataFrame(df['A'].str.split('_').to_list())
        out[2] = out[2].astype(str)
        out[4] = out[4].apply(clean_text)
        out.rename(columns={i: f'A{i}' for i in out.columns}, inplace=True)
        return pa.RecordBatch.from_pandas(out, schema=schema)
    
    此时主进程只需直接调用writer.write_batch(future.result())即可。
  • 根据你的内存情况调整chunk_size:如果内存仍有压力,可以适当减小chunk大小;如果CPU利用率不足,可以增大chunk或增加进程数。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 01:35:18