多线程读取大型CSV未占用全部CPU核心 如何优化脚本提升文件读写速度
多核CPU未占满优化方案
原代码核心问题
- 所有读写、正则匹配逻辑都在
get_next_line函数内单线程执行完毕,该函数没有返回任何可迭代对象,后续定义的4线程池全程未实际运行,自然无法占用多核CPU - 即使
get_next_line返回了行数据,现有逻辑每次处理1行就关闭、回收线程池,相当于每次仅调用1个线程处理任务,完全浪费多线程资源 - 正则匹配属于CPU密集型任务,受Python GIL(全局解释器锁)限制,多线程无法真正利用多核CPU,应改用进程池实现
- 多线程/进程直接写入同一文件会出现数据乱序、覆盖问题,原代码的写入逻辑也没有做并发兼容
优化后代码
import re import csv from datetime import datetime from multiprocessing import Pool, cpu_count from random import randint from time import sleep now = datetime.now() current_time = now.strftime("%H:%M:%S") print("Current Time start =", current_time) pattern = re.compile('some pattern') input_filename = 'in.csv' output_filename = 'out.csv' def process_line(rec): # 这里保留原有的模拟耗时逻辑,实际使用可根据需求删除或替换 # sleep(randint(0, 3)) matches = pattern.findall(rec[1]) res = [] if matches: for match in matches: res.append([rec[0], match]) return res if __name__ == '__main__': # 单线程读取所有待处理数据,避免多线程读文件产生竞争 all_recs = [] with open(input_filename, 'r', encoding='utf-8') as f: csvreader = csv.reader(f) # 若输入文件有表头需要跳过,保留下面这行,否则注释即可 # next(csvreader) all_recs = list(csvreader) # 进程数和CPU核心数保持一致,最大化利用多核性能 with Pool(processes=cpu_count()) as pool: # 批量提交所有任务并行处理 results = pool.map(process_line, all_recs) # 单线程统一写入结果,避免多进程写文件出现乱序、覆盖问题 with open(output_filename, 'w', encoding='utf-8', newline='') as csvout: csvwriter = csv.writer(csvout) csvwriter.writerow(["aa", "bb"]) for res_list in results: if res_list: csvwriter.writerows(res_list) now = datetime.now() current_time_end = now.strftime("%H:%M:%S") print("Current Time End =", current_time_end)
优化说明
- 拆分读、处理、写三个环节:读、写环节用单线程避免IO竞争,处理环节用多进程并行执行CPU密集的正则匹配任务,可完全占满4核CPU
- 移除了原代码中无效的线程池调用逻辑,批量提交所有待处理任务,最大化并行效率
- 如果你的实际业务中IO等待占比更高(比如需要调用外部接口、读写远程文件),可以将
multiprocessing.Pool改回ThreadPool,调整processes参数为你的CPU核心数即可获得最优效率
内容的提问来源于stack exchange,提问作者EagleCode
相关产品推荐
相关产品推荐

