如何在Python中并行化100GB CSV文件合并算法?
高效合并大体积CSV文件的并行解决方案
先聊聊你当前代码的问题,再给你一套可行的并行合并方案,帮你搞定这100GB的CSV合并任务。
原单线程脚本慢的核心原因
你的单线程代码有两个关键效率瓶颈:
- 每次循环都重复打开/关闭输出文件
join_rows.csv,频繁的IO操作会拖慢速度; - 逐行读取写入的方式,在Python层面做了太多循环,远不如底层的文件拷贝高效。
你的递归思路问题分析
你的分治两两合并思路是对的,但代码逻辑有不少错误:
LEN只在init=True时赋值,递归调用时不会更新,导致切片逻辑混乱;- 合并两个文件后返回的是被删除的文件路径,而不是合并后的目标文件,后续递归无法正确处理;
- 用
subprocess执行rm命令有安全风险,直接用os.remove更稳妥; - 完全没实现并行,本质还是串行处理。
可行的并行合并方案(分治+多进程)
我们用分治策略+多进程来实现并行合并,充分利用多核CPU,同时用更高效的文件拷贝方式替代逐行读写。
完整代码实现
import os import glob import re from concurrent.futures import ProcessPoolExecutor import shutil def natural_sort_key(s): # 处理文件名的自然排序(比如让file_10.csv排在file_2.csv后面) return [int(text) if text.isdigit() else text.lower() for text in re.split(r'(\d+)', s)] def merge_two_files(file1, file2, output_file, skip_header=True): """合并两个CSV文件到指定输出文件""" with open(output_file, 'wb') as out_f: # 拷贝第一个文件的全部内容 with open(file1, 'rb') as f1: shutil.copyfileobj(f1, out_f) # 拷贝第二个文件内容,跳过表头(如果需要) with open(file2, 'rb') as f2: if skip_header: next(f2) # 跳过第一行表头 shutil.copyfileobj(f2, out_f) # 删除原文件(可选,根据磁盘空间情况决定) os.remove(file1) os.remove(file2) return output_file def merge_files(file_list, skip_header=True): if len(file_list) == 1: return file_list[0] # 把文件列表分成左右两组,并行处理 mid = len(file_list) // 2 left_group = file_list[:mid] right_group = file_list[mid:] # 用多进程并行合并左右两组 with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: future_left = executor.submit(merge_files, left_group, skip_header) future_right = executor.submit(merge_files, right_group, skip_header) merged_left = future_left.result() merged_right = future_right.result() # 合并最终的两个中间文件 final_output = 'final_merged.csv' merge_two_files(merged_left, merged_right, final_output, skip_header) return final_output if __name__ == '__main__': # 替换成你的CSV文件夹路径 csv_folder = "/path/to/your/csv_files" csv_files = glob.glob(os.path.join(csv_folder, "*.csv")) # 按行位置排序文件名(自然排序) csv_files.sort(key=natural_sort_key) # 执行合并(如果你的CSV没有表头,把skip_header设为False) final_file = merge_files(csv_files, skip_header=True) print(f"合并完成!最终文件:{final_file}")
关键优化点说明
- 高效文件拷贝:用
shutil.copyfileobj替代逐行读写,这是Python调用底层系统接口的拷贝方式,速度比手动循环快得多; - 多进程并行:用
ProcessPoolExecutor实现并行分治合并,充分利用多核CPU,避免单线程的资源浪费; - 自然排序:确保文件名按行位置正确排序(比如
file_10.csv不会排在file_2.csv前面); - 表头处理:默认跳过后续文件的表头,避免合并后出现重复表头(如果你的CSV没有表头,把
skip_header设为False即可)。
单进程快速优化方案(如果不想用并行)
如果你暂时不想折腾并行,先把单线程代码优化到极致:
import glob import shutil import re def natural_sort_key(s): return [int(text) if text.isdigit() else text.lower() for text in re.split(r'(\d+)', s)] def merge_all_csvs(file_list, output_file, skip_header=True): with open(output_file, 'wb') as out_f: for idx, file_path in enumerate(file_list): with open(file_path, 'rb') as in_f: if idx > 0 and skip_header: next(in_f) # 跳过后续文件的表头 shutil.copyfileobj(in_f, out_f) out_f.write(b'\n') # 确保文件末尾有换行 # 使用示例 csv_files = sorted(glob.glob("/path/to/csvs/*.csv"), key=natural_sort_key) merge_all_csvs(csv_files, "joined.csv")
这个版本只打开一次输出文件,用高效拷贝替代逐行写入,速度会比你原来的单线程脚本快很多倍。
注意事项
- 磁盘空间:合并100GB的文件,至少需要两倍的磁盘空间(存放中间文件和最终文件),确保你的磁盘有足够余量;
- 磁盘IO瓶颈:如果是机械硬盘,不要开太多进程(比如把
max_workers设为2-4),否则磁盘IO竞争会拖慢速度; - 错误处理:实际使用时可以添加异常捕获(比如磁盘空间不足、文件读取失败等),避免中途出错前功尽弃。
内容的提问来源于stack exchange,提问作者Joe B
相关产品推荐
相关产品推荐

