PyQt5非主线程启动多进程实现Excel数据逐行API发送方案问询
实现方案
核心实现逻辑要先解决几个关键问题:
- openpyxl的工作表对象无法直接跨进程序列化传输,所以需要提前将所有待处理行的所需数据全部读取为可序列化的字典/列表,再提交给进程池处理
- 要限制最大并发进程数,避免短时间请求量过大触发接口限流,也避免占用过多系统资源
- 所有统计变量统一在WorkerThread中汇总,每个进程只负责处理单行数据并返回处理结果
修改后的核心代码
改造api_import函数
去掉对sheet对象的依赖,入参改为预读取的可序列化行数据:
def api_import(row_data): # row_data结构:{'item': 行item值, 'templ_id': 模板id, 'params': [所有参数字典]} item = row_data['item'] templ_id = row_data['templ_id'] params = row_data['params'] item_id = get_item_id(item)['item_id'] query = {"template_id": str(templ_id)} resp = templ_attach(item_id, query) ok = 1 if resp == 200 else 0 item_error = 0 for param in params: param_name = param['name'] param_description = param['description'] param_value = param['value'] param_type = param['type'] if "[T]" in param_type: query_param = {"name": str(param_name), "type": 4, "description": str(param_description), "text_value": str(param_value)} elif "[D]" in param_type: query_param = {"name": str(param_name), "type": 0, "description": str(param_description), "numeric_value": param_value} else: query_param = {"name": str(param_name), "type": 1, "description": str(param_description), "logical_value": param_value} resp = input_param(item_id, query_param) if resp not in (201, 202): item_error += 1 total_success = 1 if item_error == 0 else 0 return {'ok': ok, 'total_success': total_success}
改造WorkerThread的run方法
提前读取所有Excel数据构造任务列表,用进程池并发处理:
import time from concurrent.futures import ProcessPoolExecutor, as_completed from openpyxl import load_workbook from openpyxl.utils import get_column_letter class WorkerThread(QThread): update_progress = pyqtSignal(int) worker_complete = pyqtSignal(dict) def run(self): start_time = time.time() # 可改为从主窗口获取用户选择的文件路径,替换原写死的路径 file_path = "import.xlsx" wb = load_workbook(file_path, read_only=True) # 只读模式降低大文件内存占用 all_tasks = [] items_number = 0 # 提前读取所有待处理数据,构造可序列化的任务列表 for sheet in wb: number_rows = sheet.max_row number_columns = sheet.max_column items_number += (number_rows - 3) # 先读取表头通用配置 header_info = {} for l in range(4, number_columns + 1): col_letter = get_column_letter(l) header_info[l] = { 'name': sheet[col_letter + '1'].value.upper().replace(" ", ""), 'description': sheet[col_letter + '2'].value, 'type': sheet[col_letter + '3'].value } # 读取每行业务数据 for k in range(4, number_rows + 1): row_data = { 'item': sheet[get_column_letter(3) + str(k)].value, 'templ_id': sheet[get_column_letter(number_columns) + str(k)].value, 'params': [] } for l in range(4, number_columns + 1): param = header_info[l].copy() param['value'] = sheet[get_column_letter(l) + str(k)].value row_data['params'].append(param) all_tasks.append(row_data) wb.close() # 多进程处理任务,max_workers可根据接口承压能力调整,建议不要超过8 ok = 0 total_success = 0 finished = 0 with ProcessPoolExecutor(max_workers=4) as executor: futures = [executor.submit(api_import, task) for task in all_tasks] for future in as_completed(futures): res = future.result() ok += res['ok'] total_success += res['total_success'] finished += 1 progress = round(finished / items_number * 100) self.update_progress.emit(progress) # 完成后返回全量统计数据 self.worker_complete.emit({ "emp_id":1234, "fn":"XXX", "ln":"YYYY", "total": items_number, "ok": ok, "total_success": total_success }) end_time = time.time() - start_time item_time = end_time / items_number
额外注意事项
- Windows环境下运行代码必须将主程序启动逻辑包裹在
if __name__ == '__main__':块中,否则多进程会重复启动GUI导致报错 - 如果接口有请求频率限制,可以调小max_workers参数,或者在api_import中添加合理的休眠逻辑
- 可以额外添加异常捕获逻辑,处理API请求超时、网络错误等异常情况,避免进程崩溃
- 如果需要从主窗口传递文件路径,可以给WorkerThread加初始化参数传递选好的文件地址
内容的提问来源于stack exchange,提问作者Michał Sarna
相关产品推荐
相关产品推荐

