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

多CSV文件并行处理:如何从multiprocessing函数返回counter字典?

并行处理CSV文件并收集Counter字典的解决方案

嘿,刚好我之前处理过几乎一模一样的需求,结合你提到的那个进程返回值方案,给你整理了一个直接能用的实现,完美解决并行处理CSV并收集Counter字典的问题👇

核心思路

用multiprocessing.Process配合**队列(Queue)**来收集每个子进程的返回结果——这也是你提到的方案的核心逻辑:让每个子进程处理完CSV后,把生成的Counter字典放进队列,主进程等所有子进程跑完后,统一从队列里取出结果拼成主列表。

完整代码实现

1. 导入依赖模块

import csv
from collections import Counter
import multiprocessing

2. 定义单个CSV的处理函数

这个函数负责读取单个CSV、生成Counter字典,并把结果存入队列:

def process_csv(file_path, result_queue):
    # 初始化空Counter
    counter = Counter()
    
    # 读取CSV文件(这里假设统计第0列内容,可根据你的需求修改)
    with open(file_path, 'r', encoding='utf-8') as f:
        reader = csv.reader(f)
        # 跳过表头(不需要的话直接删掉这行)
        next(reader)
        for row in reader:
            # 统计目标列,这里以第0列为例,按需调整
            target_item = row[0]
            counter[target_item] += 1
    
    # 将处理好的Counter存入队列,供主进程收集
    result_queue.put(counter)

3. 主函数:管理进程并汇总结果

def main(csv_files):
    # 创建用于存储结果的队列
    result_queue = multiprocessing.Queue()
    processes = []
    
    # 为每个CSV文件创建一个子进程
    for file in csv_files:
        process = multiprocessing.Process(
            target=process_csv,
            args=(file, result_queue)
        )
        processes.append(process)
        process.start()
    
    # 等待所有子进程执行完毕
    for process in processes:
        process.join()
    
    # 从队列中取出所有Counter,组成主列表
    main_counter_list = []
    while not result_queue.empty():
        main_counter_list.append(result_queue.get())
    
    return main_counter_list

4. 调用示例

if __name__ == "__main__":
    # 替换成你的CSV文件路径列表
    my_csv_files = ["data1.csv", "data2.csv", "data3.csv"]
    
    # 执行并行处理
    final_result_list = main(my_csv_files)
    
    # 打印结果(可根据需求调整输出逻辑)
    for idx, counter in enumerate(final_result_list):
        print(f"=== 文件 {my_csv_files[idx]} 的统计结果 ===")
        print(counter)

可选简化方案:用multiprocessing.Pool

如果你不需要精细控制进程(比如设置进程名、优先级),用Pool.map会更简洁,它会自动帮你管理进程池并返回结果列表:

import csv
from collections import Counter
from multiprocessing import Pool

def process_csv(file_path):
    counter = Counter()
    with open(file_path, 'r', encoding='utf-8') as f:
        reader = csv.reader(f)
        next(reader)
        for row in reader:
            counter[row[0]] += 1
    return counter

if __name__ == "__main__":
    my_csv_files = ["data1.csv", "data2.csv", "data3.csv"]
    
    # 用CPU核心数作为进程池大小,也可以手动指定
    with Pool(processes=multiprocessing.cpu_count()) as pool:
        final_result_list = pool.map(process_csv, my_csv_files)
    
    print(final_result_list)

关键注意事项

  • Windows系统兼容性:必须把主程序入口放在if __name__ == "__main__"代码块里,否则会出现进程重复启动的问题,这是Python多进程在Windows下的特殊要求。
  • 自定义统计逻辑:如果需要统计CSV的其他内容(比如多列、特定条件的行),只需要修改process_csv函数里的业务逻辑即可。
  • 队列的安全性:用队列传递结果是多进程间安全的通信方式,不会出现共享变量的竞争问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:54:15