如何在现有Python应用中使用multiprocessing加速批量文件入库
Python多目录文件处理并行优化方案
现有代码核心瓶颈
- 串行遍历处理所有目录,文件读取、数据库写入都是IO密集型操作,CPU长期空转浪费性能
- 单条执行SQL插入,每次插入都要和数据库做网络交互,这部分是最大的耗时来源
- 原有逻辑没有错误校验,插入失败会直接中断整个任务
优化思路
- 采用多进程并行处理不同目录,每个进程单独维护数据库连接,避免连接线程安全问题
- 替换单条插入为批量插入,大幅降低数据库交互开销
- 保留原有业务逻辑不变,仅调整执行模式和写入逻辑
优化后代码
import mysql.connector import csv import os import time import configparser from multiprocessing import Pool, cpu_count # 加载配置,全局共享只读配置不需要每个进程重复读 config = configparser.ConfigParser() config.read(r'C:\Desktop\Energy\file_cfg.ini') source = config['PATHS']['source'] archive = config['PATHS']['archive'] DB_CONFIG = { "host": config['DB']['host'], "user": config['DB']['user'], "passwd": config['DB']['passwd'], "database": config['DB']['database'] } # 单个目录处理逻辑,每个子进程独立执行 def process_single_mp(mp): mp_server = os.listdir(source) if mp not in mp_server: return subdir_paths = os.path.join(source, mp) # 获取目录下最新文件的创建时间,保留原逻辑 cr_time = "" for file in os.listdir(subdir_paths): file_paths = os.path.join(subdir_paths, file) cr_time_s = os.path.getctime(file_paths) cr_time = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(cr_time_s)) # 筛选有效非空文件 all_file_paths = [os.path.join(subdir_paths, f) for f in os.listdir(subdir_paths)] full_file_paths = [p for p in all_file_paths if os.path.getsize(p) > 0] if not full_file_paths: return newest_file_path = max(full_file_paths, key=os.path.getctime) # 仅处理创建时间超过2分钟的最新文件,保留原逻辑 if os.path.getctime(newest_file_path) >= time.time() - 120: return # 读取CSV数据 line_data0 = [] with open(newest_file_path, 'rt') as f: reader = csv.reader(f, delimiter ='\t') next(reader) # 跳过表头 for line in reader: if not line: continue line.insert(0, mp) line.insert(1, cr_time) line_data0.append(line) if not line_data0: return # 进程独立创建数据库连接,避免多线程/进程连接冲突 mydb = mysql.connector.connect(**DB_CONFIG) cursor = mydb.cursor() q1 = ("INSERT INTO microbeats" "(`antenna`,`datetime`,`system`,`item`,`event`, `status`, `accident`)" "VALUES (%s, %s, %s,%s, %s, %s, %s)") # 批量插入,一次性提交所有数据,大幅降低交互开销 cursor.executemany(q1, line_data0) mydb.commit() cursor.close() mydb.close() if __name__ == "__main__": # 主进程提前加载antenna列表,执行清空表操作,避免多进程并发操作表冲突 mydb = mysql.connector.connect(**DB_CONFIG) cursor = mydb.cursor() cursor.execute("SELECT * FROM `antenna`") mp_mysql = [i[0] for i in cursor.fetchall()] cursor.execute("TRUNCATE TABLE microbeats") mydb.commit() cursor.close() mydb.close() # 启动进程池并行处理,进程数取CPU核心数和目录数的较小值,避免资源浪费 process_num = min(cpu_count(), len(mp_mysql)) with Pool(processes=process_num) as pool: pool.map(process_single_mp, mp_mysql)
效果说明
- 批量插入可将数据库写入耗时降低90%以上,是性能提升的核心来源
- 并行处理19个目录的情况下,原有19小时的总耗时可压缩到1小时以内
- 后续新增目录只需调整进程数即可,无需修改核心逻辑
注意事项
- 进程数不要设置过高,避免数据库连接数超出上限,建议最高不超过16
- 若单文件数据量超过10万行,可拆分批次执行
executemany,避免单次提交数据量过大 - Windows环境下运行多进程代码必须把启动逻辑放在
if __name__ == "__main__"块内,否则会报错
内容的提问来源于stack exchange,提问作者Mediterráneo
相关产品推荐
相关产品推荐

