ThreadPoolExecutor多用户场景卡顿问题求助
开发了一个Python机器人,接收用户指令后将多个ID发送至API,存储结果返回给用户。为提升效率,采用ThreadPoolExecutor,每个用户设置max_workers=5。1-2个用户使用正常,但50-80个用户同时触发指令时,用户仅收到“Checking...”提示,后续无响应。调用的API无问题,每秒1万次请求也能在1-1.5秒内返回结果,CPU使用率仅25-30%,推测线程处于等待状态。
代码如下:
from pyrogram import Client, filters import requests import concurrent.futures import threading import asyncio import concurrent.futures def do_work(id): while True: result = requests.get(f"example.com/api.php?id={id}") if result.status_code == 200: return result.text else: continue @Client.on_message(filters.command("getresult", [".", "/"])) def multi(Client, message): t1 = threading.Thread(target=bcall, args=(Client, message)) t1.start() def bcall(Client, message): loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) loop.run_until_complete(thread(Client, message)) loop.close() async def thread(Client, message): ids = message.text.split("\n") ALL_RESULT = "" await message.reply_text("Checking...") with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor: futures = [] for i in ids: future = executor.submit(do_work , i) futures.append(future) for future in concurrent.futures.as_completed(futures): result = future.result() ALL_RESULT += result + "\n" await message.reply_text(ALL_RESULT) await message.reply_text(f"Done Checking {len(ids)} IDS")
注:do_work中的while循环实际已设置max_retry=2后终止,会在2-3秒内结束。
1. 移除冗余线程嵌套,利用Pyrogram异步特性
当前代码在同步函数里开线程,线程内再手动创建asyncio事件循环,这种嵌套会快速耗尽线程资源。Pyrogram本身是异步框架,直接用异步函数处理即可:
@Client.on_message(filters.command("getresult", [".", "/"])) async def multi(client, message): await thread(client, message)
直接删除bcall函数,无需手动管理事件循环。
2. 替换同步requests为异步HTTP库
requests是同步阻塞库,大量并发时会因IO阻塞占用线程,导致线程池耗尽。改用aiohttp实现异步请求,配合asyncio.gather实现高效并发:
先安装依赖:pip install aiohttp
修改do_work为异步函数:
import aiohttp async def do_work(session, id): max_retry = 2 for _ in range(max_retry + 1): try: async with session.get(f"https://example.com/api.php?id={id}") as response: if response.status == 200: return await response.text() except Exception: pass return f"Failed to get result for ID {id}"
更新thread函数,用异步并发替代线程池:
async def thread(client, message): # 过滤空ID,避免无效请求 ids = [id.strip() for id in message.text.split("\n") if id.strip()] await message.reply_text("Checking...") async with aiohttp.ClientSession() as session: tasks = [do_work(session, id) for id in ids] results = await asyncio.gather(*tasks) # 用join拼接结果,比循环+=更高效 all_result = "\n".join(results) await message.reply_text(all_result) await message.reply_text(f"Done Checking {len(ids)} IDS")
3. 限制全局并发请求量
50-80个用户同时请求时,总并发数可能过高。用asyncio.Semaphore控制全局并发上限,避免资源耗尽:
# 在文件顶部定义全局信号量 MAX_CONCURRENT_REQUESTS = 50 semaphore = asyncio.Semaphore(MAX_CONCURRENT_REQUESTS) async def do_work(session, id): max_retry = 2 for _ in range(max_retry + 1): async with semaphore: try: async with session.get(f"https://example.com/api.php?id={id}") as response: if response.status == 200: return await response.text() except Exception: pass return f"Failed to get result for ID {id}"
4. 优化结果拼接逻辑
原代码用ALL_RESULT += result + "\n"做循环拼接,字符串是不可变对象,频繁拼接会产生大量中间对象,内存开销大。改用列表收集结果后用"\n".join(results)拼接,效率更高。
5. 规避Telegram API限制
大量并发回复消息可能触发Telegram的频率限制,可给回复操作加短暂延迟,或对过长结果拆分后发送:
# 示例:给每条回复加0.1秒延迟 await asyncio.sleep(0.1) await message.reply_text(all_result)
内容的提问来源于stack exchange,提问作者Mainul

