如何高效将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}")
关键改动说明
- 用
as_completed替代map:不再等待所有任务完成,而是实时获取已完成的转换结果 - 主进程负责写入:ParquetWriter由主进程创建和调用,避免多进程共享资源的安全问题
- 无合并操作:每个转换后的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
相关产品推荐
相关产品推荐

