Python多进程/多线程新手咨询:任务并行方案选型
多线程/多进程优化员工测试分析项目的问题解答
背景概述
接触Python约一周,为提升程序速度开展测试项目,学习多线程与多进程。现有5000条员工记录需分析、50个自定义Python测试、结果存储表,运行环境为Windows 4核i7处理器,流程为:
- a) 遍历员工记录
- b) 遍历测试
- c) 运行测试
- d) 记录测试结果
针对提出的四个问题,逐一解答如下:
1. 步骤(a)读取员工记录是否用多线程?
不需要。数据库I/O的优化核心是减少连接次数和批量操作,而非多线程并发读取。5000条员工ID完全可以用单线程一次性从数据库批量查询出来存入内存列表,后续直接遍历内存中的数据即可。多线程读取反而会因为数据库连接竞争、锁机制等问题降低效率,甚至引发连接超时。
2. 步骤(b)读取测试记录是否用多线程?
完全没必要。50个测试的代码量极小,单线程一次性把所有TestCode查询出来存入字典缓存(比如{test_id: test_code}),后续所有测试任务直接从内存缓存取代码,彻底避免重复的数据库I/O。多线程在这里不仅无收益,还会增加代码复杂度。
3. 步骤(c)运行测试是否用多进程?
必须用多进程。测试运行属于CPU密集型操作,Python的GIL(全局解释器锁)会限制多线程无法利用多核CPU,而多进程可以绕过GIL,充分发挥4核i7的性能。但注意进程数量不要超过CPU核心数的2倍(比如4-8个),过多进程会导致频繁的进程切换,反而拖慢整体速度。
4. 步骤(d)写入测试结果是否用多线程?
不建议直接对每个写入操作开多线程。数据库写入的瓶颈在于磁盘I/O和锁机制,多线程并发写入容易引发锁竞争,导致commit等待时间变长。更优的方案是:
- 把所有测试结果先存入一个线程安全的队列(比如
queue.Queue) - 启动1-2个专门的写入线程,从队列中批量取出结果,一次性插入数据库并批量commit,大幅减少I/O次数。
当前代码的关键问题
- 多进程共享SQLite连接:全局的
connection和cursor不能在多进程间共享,每个进程必须创建独立的数据库连接,否则会出现数据损坏、连接报错等问题。 - 进程数量过多:当前代码会启动50个员工进程,每个又启动10个测试进程,总共500个进程,远超过4核CPU的处理能力,会导致严重的进程切换开销。
- SQL注入风险:
GetCode函数中用字符串格式化拼接SQL语句,存在SQL注入风险,应使用参数化查询。 - 语法错误:
testDate = datetime.datetime.now缺少括号,应为testDate = datetime.datetime.now()。 - 频繁commit:每个测试都单独
commit,会大幅增加数据库I/O开销,应批量提交。
优化后的代码示例
import datetime from multiprocessing import Pool, Manager import sqlite3 import threading # 全局缓存测试代码,单线程提前加载 TEST_CODE_CACHE = {} def load_test_codes(): """提前加载所有测试代码到缓存""" conn = sqlite3.connect('EmpTestPython.db', timeout=90) cursor = conn.cursor() cursor.execute('SELECT ID, TestCode FROM Tests') for test_id, test_code in cursor.fetchall(): TEST_CODE_CACHE[test_id] = test_code conn.close() def run_single_test(args): """单个测试任务:运行测试并返回结果(不直接写入DB)""" emp_id, test_id = args test_code = TEST_CODE_CACHE[test_id] test_date = datetime.datetime.now() # 模拟测试运行,实际场景中eval有安全风险,建议替换为更安全的执行方式 result = eval(test_code) return (test_date, emp_id, test_id, result) def write_results_to_db(result_queue): """专门的写入线程:批量处理结果写入""" conn = sqlite3.connect('EmpTestPython.db', timeout=90) cursor = conn.cursor() batch = [] while True: result = result_queue.get() if result is None: # 结束信号 break batch.append(result) # 每100条批量提交一次,可根据实际情况调整批次大小 if len(batch) >= 100: cursor.executemany(''' INSERT INTO TestResults (TestDate, EmpId, TestId, Result) VALUES (?, ?, ?, ?) ''', batch) conn.commit() batch = [] # 提交剩余的结果 if batch: cursor.executemany(''' INSERT INTO TestResults (TestDate, EmpId, TestId, Result) VALUES (?, ?, ?, ?) ''', batch) conn.commit() conn.close() if __name__ == '__main__': # 1. 提前加载测试代码缓存 load_test_codes() # 2. 批量读取所有员工ID conn = sqlite3.connect('EmpTestPython.db', timeout=90) cursor = conn.cursor() cursor.execute('SELECT ID FROM Employees') emp_ids = [row[0] for row in cursor.fetchall()] conn.close() # 3. 生成所有测试任务:每个员工对应50个测试 tasks = [] for emp_id in emp_ids: for test_id in TEST_CODE_CACHE.keys(): tasks.append((emp_id, test_id)) # 4. 启动结果队列和写入线程 manager = Manager() result_queue = manager.Queue() write_thread = threading.Thread(target=write_results_to_db, args=(result_queue,)) write_thread.start() # 5. 用进程池运行测试任务,进程数设为CPU核心数(4核设为4) with Pool(processes=4) as pool: for result in pool.imap_unordered(run_single_test, tasks): result_queue.put(result) # 6. 发送结束信号,等待写入线程完成 result_queue.put(None) write_thread.join() print("所有测试完成,结果已写入数据库")
内容的提问来源于stack exchange,提问作者lmoody
相关产品推荐
相关产品推荐

