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

如何为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])])

为啥之前并行失败?大概率是这几个坑

  1. 没加if __name__ == '__main__'::Windows系统下多进程必须加这个,不然会重复执行代码导致混乱。
  2. 共享文件对象:直接在多进程里用同一个文件指针,会导致读取位置错乱。
  3. 对比逻辑有状态:比如依赖全局变量,多进程下每个进程的全局变量是独立的,会导致结果不对。

额外小优化

  • 如果你的对比是IO密集型(比如还要读写其他文件),可以把ProcessPoolExecutor换成ThreadPoolExecutor,线程开销更小。
  • 先检查两个CSV的行数是否一致,避免出现行不匹配的情况。
  • 要是CSV有表头,用csv.DictReader按列名对比,代码可读性会更好。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:31:30