如何在现有协程中提升asyncio的线程限制?
线程被限制的原因
asyncio.to_thread 默认复用 asyncio 事件循环绑定的全局线程池,这个线程池的max_workers默认值为20(对应Python标准库concurrent.futures.ThreadPoolExecutor的默认配置)。当提交的阻塞任务数量超过20时,多余任务会进入线程池等待队列,只能等已有线程空闲后才能执行——这就是为什么任务数超过20时,asyncio.gather的执行时间会翻倍。
哪怕你增加worker协程的数量也没用,因为所有asyncio.to_thread调用都共享这个全局线程池,线程总数被池的上限直接卡死。
提升线程限制的方法(沿用asyncio方案)
不用全量转异步或切换多进程,有两种简单方式调整线程池大小:
1. 替换全局默认线程池
在程序启动时,给事件循环设置自定义全局线程池,指定更大的max_workers,后续所有asyncio.to_thread调用都会自动使用这个池:
import asyncio from concurrent.futures import ThreadPoolExecutor async def main(): loop = asyncio.get_running_loop() # 设置全局线程池最大线程数为100 loop.set_default_executor(ThreadPoolExecutor(max_workers=100)) # 批量执行阻塞任务 await asyncio.gather( *[asyncio.to_thread(lambda: time.sleep(1)) for _ in range(61)] ) if __name__ == "__main__": asyncio.run(main())
2. 为特定任务使用独立线程池
如果不想修改全局配置,可以手动创建线程池,在调用asyncio.to_thread时通过executor参数指定使用该池,灵活性更高:
import asyncio from concurrent.futures import ThreadPoolExecutor import time # 模拟阻塞逻辑的同步函数 def blocking_task(task_id): time.sleep(1) print(f"任务 {task_id} 完成") async def worker(queue, executor): while True: task_id = await queue.get() await asyncio.to_thread(blocking_task, task_id, executor=executor) queue.task_done() async def operation_loop(): # 创建自定义线程池,设置最大线程数为60 custom_executor = ThreadPoolExecutor(max_workers=60) queue = asyncio.Queue() # 启动30个worker协程 workers = [asyncio.create_task(worker(queue, custom_executor)) for _ in range(30)] # 每秒填充任务,模拟业务场景 for i in range(61): await queue.put(i) await asyncio.sleep(1) await queue.join() # 清理资源 for w in workers: w.cancel() await asyncio.gather(*workers, return_exceptions=True) custom_executor.shutdown(wait=True) if __name__ == "__main__": asyncio.run(operation_loop())
注意:调整线程数要结合实际场景(如数据库连接上限、系统CPU/内存资源)合理设置,并非越大越好——过多线程会增加上下文切换开销,反而降低整体性能。
内容的提问来源于stack exchange,提问作者Simarra
相关产品推荐
相关产品推荐

