如何为Python大CSV文件对比代码实现CPU全核并行处理?
搞定CSV并行对比!让8核CPU跑满的方案
嘿,我明白你的痛点——10万行的CSV单线程跑太慢,8核CPU闲着简直浪费!之前并行尝试失败大概率是踩了Python多进程的几个常见坑,我给你一套直接能用的方案,分两种场景(内存足够/内存吃紧),一步步来:
核心思路
咱们把大文件拆成8块(对应你的8核),让每个核心处理一块数据的对比,最后把所有结果合并写入输出文件。Python里concurrent.futures.ProcessPoolExecutor是最省心的工具,API简单还能自动管理进程。
先搞清楚几个关键注意点
- 多进程里不能直接共享文件对象!每个进程得自己读对应的内容,不然会搞乱文件指针。
- 你的对比逻辑得是独立无状态的——就是每行的对比结果不依赖其他行,这样拆分才安全。
场景1:内存足够(推荐!10万行CSV占内存很小)
先把两个CSV全读进内存,这样后续并行时不用反复读文件,效率最高:
import csv from concurrent.futures import ProcessPoolExecutor # 先写个加载CSV的小函数 def load_csv(file_path): with open(file_path, 'r', newline='', encoding='utf-8') as f: return list(csv.reader(f)) # 把两个CSV加载成列表 host_data = load_csv('host.csv') master_data = load_csv('master.csv')
然后写你的对比逻辑(这里假设你是逐行对比所有列,你可以改成自己的需求):
def compare_single_row(args): row_idx, host_row, master_row = args # 这里写你的对比规则,比如检查列是否相等 if host_row != master_row: # 返回差异行的索引和内容,方便后续写入 return (row_idx + 1, host_row, master_row) # 行号从1开始更友好 return None
最后启动进程池跑起来:
if __name__ == '__main__': # 准备任务:把两个CSV的行一一配对,带上索引 tasks = [(i, host_data[i], master_data[i]) for i in range(len(host_data))] # 开8个进程(和CPU核心数一致) with ProcessPoolExecutor(max_workers=8) as executor: # 批量提交任务,获取结果 all_results = executor.map(compare_single_row, tasks) # 把差异结果写入输出文件 with open('results.csv', 'w', newline='', encoding='utf-8') as result_file: writer = csv.writer(result_file) writer.writerow(['行号', 'host.csv内容', 'master.csv内容']) # 加个表头 for res in all_results: if res is not None: idx, h_row, m_row = res writer.writerow([idx, ','.join(h_row), ','.join(m_row)])
场景2:内存吃紧(比如文件更大)
如果内存不够加载整个文件,就用分块读取的方式,每个进程处理一块:
import csv from concurrent.futures import ProcessPoolExecutor def process_chunk(chunk_info): start_row, end_row, host_path, master_path = chunk_info differences = [] # 处理当前块的每一行 for row_idx in range(start_row, end_row): # 读取host.csv的对应行 with open(host_path, 'r', newline='', encoding='utf-8') as host_file: host_reader = csv.reader(host_file) # 跳过前面的行 for _ in range(row_idx): next(host_reader, None) host_row = next(host_reader, None) # 读取master.csv的对应行 with open(master_path, 'r', newline='', encoding='utf-8') as master_file: master_reader = csv.reader(master_file) for _ in range(row_idx): next(master_reader, None) master_row = next(master_reader, None) # 对比逻辑,和之前一样 if host_row != master_row and host_row is not None and master_row is not None: differences.append((row_idx + 1, host_row, master_row)) return differences if __name__ == '__main__': total_rows = 100000 chunk_size = total_rows // 8 # 每块约12500行 chunks = [] # 拆分出8个块,最后一块处理剩余的行 for i in range(8): start = i * chunk_size end = start + chunk_size if i < 7 else total_rows chunks.append((start, end, 'host.csv', 'master.csv')) # 启动进程池处理所有块 with ProcessPoolExecutor(max_workers=8) as executor: chunk_results = executor.map(process_chunk, chunks) # 合并所有块的结果写入文件 with open('results.csv', 'w', newline='', encoding='utf-8') as result_file: writer = csv.writer(result_file) writer.writerow(['行号', 'host.csv内容', 'master.csv内容']) for chunk_res in chunk_results: for diff_row in chunk_res: writer.writerow([diff_row[0], ','.join(diff_row[1]), ','.join(diff_row[2])])
为啥之前并行失败?大概率是这几个坑
- 没加
if __name__ == '__main__'::Windows系统下多进程必须加这个,不然会重复执行代码导致混乱。 - 共享文件对象:直接在多进程里用同一个文件指针,会导致读取位置错乱。
- 对比逻辑有状态:比如依赖全局变量,多进程下每个进程的全局变量是独立的,会导致结果不对。
额外小优化
- 如果你的对比是IO密集型(比如还要读写其他文件),可以把
ProcessPoolExecutor换成ThreadPoolExecutor,线程开销更小。 - 先检查两个CSV的行数是否一致,避免出现行不匹配的情况。
- 要是CSV有表头,用
csv.DictReader按列名对比,代码可读性会更好。
内容的提问来源于stack exchange,提问作者Yoekleng Kuy
相关产品推荐
相关产品推荐

