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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 19:35:36