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

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次数。

当前代码的关键问题

  1. 多进程共享SQLite连接:全局的connection和cursor不能在多进程间共享,每个进程必须创建独立的数据库连接,否则会出现数据损坏、连接报错等问题。
  2. 进程数量过多:当前代码会启动50个员工进程,每个又启动10个测试进程,总共500个进程,远超过4核CPU的处理能力,会导致严重的进程切换开销。
  3. SQL注入风险:GetCode函数中用字符串格式化拼接SQL语句,存在SQL注入风险,应使用参数化查询。
  4. 语法错误:testDate = datetime.datetime.now缺少括号,应为testDate = datetime.datetime.now()。
  5. 频繁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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 16:43:31