Python ThreadPool实现带超时的持续并发负载测试及结果捕获
负载测试脚本解决方案(维持固定并行线程+超时控制+全结果捕获)
问题分析
原代码存在两个核心问题:
- 一次性提交所有任务,无法维持固定数量的并行线程直到超时(任务完成后不会自动补充新任务)
- 没有超时机制,也无法在超时后捕获所有线程的运行状态
解决方案代码
import datetime import time import threading from concurrent.futures import ThreadPoolExecutor, wait, FIRST_COMPLETED, TimeoutError def query_data(q, stop_event): start_time = datetime.datetime.utcnow() task_id = q[0] # 模拟业务逻辑(拆分长sleep为多次短sleep,以便响应停止信号) for _ in range(50): if stop_event.is_set(): end_time = datetime.datetime.utcnow() run_time = (end_time - start_time).total_seconds() return {"id": task_id, "run_time": round(run_time, 2), "status": "cancelled"} time.sleep(0.1) end_time = datetime.datetime.utcnow() run_time = (end_time - start_time).total_seconds() return {"id": task_id, "run_time": round(run_time, 2), "status": "success"} def run_processes(queries, max_threads=5, timeout=30): stop_event = threading.Event() results = [] start_time = datetime.datetime.utcnow() with ThreadPoolExecutor(max_workers=max_threads) as executor: # 初始化线程池:先提交max_threads个任务 futures = {} for q in queries[:max_threads]: future = executor.submit(query_data, q, stop_event) futures[future] = q remaining_queries = queries[max_threads:] # 循环维持线程负载,直到超时 while futures and (datetime.datetime.utcnow() - start_time).total_seconds() < timeout: # 等待1秒内有任务完成,避免主线程阻塞 done, not_done = wait(futures, timeout=1, return_when=FIRST_COMPLETED) # 处理已完成的任务,补充新任务 for future in done: try: result = future.result() results.append(result) except Exception as e: task_id = futures[future][0] results.append({"id": task_id, "status": "failed", "error": str(e)}) # 有剩余任务就提交新的,保持线程池满负载 if remaining_queries: q = remaining_queries.pop(0) new_future = executor.submit(query_data, q, stop_event) futures[new_future] = q del futures[future] # 触发停止信号,通知所有运行中的任务主动退出 stop_event.set() # 处理剩余未完成的任务,收集最终状态 for future in futures: task_id = futures[future][0] try: # 给任务5秒时间主动退出 result = future.result(timeout=5) results.append(result) except TimeoutError: results.append({"id": task_id, "status": "force_terminated"}) except Exception as e: results.append({"id": task_id, "status": "failed", "error": str(e)}) return results # 测试示例 if __name__ == "__main__": test_queries = [(i,) for i in range(20)] test_results = run_processes(test_queries, max_threads=5, timeout=10) for res in test_results: print(res)
关键特性说明
- 固定并行负载维持:初始提交
max_threads个任务,每当任务完成立即补充新任务,始终保持指定数量的线程并行运行。 - 超时控制:主线程实时监控运行时间,达到超时阈值后停止提交新任务,并通知正在运行的任务主动退出。
- 全结果捕获:所有任务的状态(成功、取消、失败、强制终止)都被记录,包括超时后未完成的任务信息。
- 协作式线程退出:通过
stop_event让线程主动响应停止信号,规避Python无法强制终止线程的限制。
内容的提问来源于stack exchange,提问作者Incognito
相关产品推荐
相关产品推荐

