如何用Asyncio在单个Python应用中并发运行多个类实例?
我尝试在单个Python应用中并发运行同一个类的多个实例,模拟每个实例在独立终端运行的效果,目标是让每个实例独立运行且互不阻塞,目前使用asyncio实现并发。理论上它们确实在并发运行,但某个Operations实例中的await调用似乎会导致其他所有Operations实例也等待。
简化实现代码
import os import asyncio # 原代码遗漏asyncio导入,补充后保证可运行 # 加载环境变量 os.environ['ADDRESS1'] = 'address1' os.environ['PRIVATE_KEY1'] = 'private_key1' os.environ['ADDRESS2'] = 'address2' os.environ['PRIVATE_KEY2'] = 'private_key2' class Operations: def __init__(self, address, private_key): self.address = address self.private_key = private_key async def run(self): while True: print(f"Running operations for {self.address}") await do_other_thing() await asyncio.sleep(1) # 假设这是自定义业务函数,可能存在同步阻塞问题 def do_other_thing(): # 模拟同步阻塞操作(如同步网络请求、本地IO、CPU计算) import time time.sleep(2) # 同步sleep会阻塞整个asyncio事件循环 # 初始化不同配置的实例 instances = [ Operations(os.getenv('ADDRESS1'), os.getenv('PRIVATE_KEY1')), Operations(os.getenv('ADDRESS2'), os.getenv('PRIVATE_KEY2')) ] async def run_instance(instance): await instance.run() async def main(): tasks = [asyncio.create_task(run_instance(instance)) for instance in instances] await asyncio.gather(*tasks) if __name__ == "__main__": asyncio.run(main())
问题核心
尽管使用asyncio.create_task并发启动实例,但每个Operations实例中的await调用会导致全局延迟,实例无法真正独立运行。
已尝试方案
- 确保
run方法内所有操作都使用await,预留任务切换时机; - 添加短暂sleep避免忙循环。
需求
指导如何在单个Python应用中高效并发运行多个类实例,让它们的表现如同在独立终端中运行一样,最大化并发能力且互不阻塞。
解决方案
1. 排查核心阻塞点:do_other_thing()
你的问题大概率出在do_other_thing()上:如果这个函数是同步阻塞类型(包含time.sleep()、同步IO、CPU密集计算),即使加了await也会阻塞整个asyncio事件循环——asyncio的并发依赖非阻塞异步操作,只有当任务遇到真正的IO等待时,事件循环才会切换到其他任务。同步代码会霸占事件循环,导致所有实例等待。
2. 用线程池包装同步阻塞操作
如果do_other_thing()无法改造成异步函数(比如依赖第三方同步库),可以用asyncio.to_thread()(Python3.9+支持)将其放到线程池运行,避免阻塞事件循环:
修改后的run方法示例:
async def run(self): while True: print(f"Running operations for {self.address}") # 将同步函数放到独立线程执行 await asyncio.to_thread(do_other_thing) await asyncio.sleep(1)
也可以手动使用ThreadPoolExecutor控制线程数量:
from concurrent.futures import ThreadPoolExecutor # 初始化线程池,根据任务数量调整max_workers executor = ThreadPoolExecutor(max_workers=4) async def run(self): while True: print(f"Running operations for {self.address}") await asyncio.get_running_loop().run_in_executor(executor, do_other_thing) await asyncio.sleep(1)
3. CPU密集型任务用多进程
如果任务是CPU密集型(比如大量数值计算),asyncio的单线程模型无法利用多核,此时应该用multiprocessing让每个实例在独立进程中运行:
示例代码:
import os import time import multiprocessing os.environ['ADDRESS1'] = 'address1' os.environ['PRIVATE_KEY1'] = 'private_key1' os.environ['ADDRESS2'] = 'address2' os.environ['PRIVATE_KEY2'] = 'private_key2' class Operations: def __init__(self, address, private_key): self.address = address self.private_key = private_key def run(self): # 用普通同步方法即可 while True: print(f"Running operations for {self.address} (进程ID: {os.getpid()})") do_other_thing() time.sleep(1) def do_other_thing(): # 模拟CPU密集型计算 sum(range(10**6)) if __name__ == "__main__": instances = [ Operations(os.getenv('ADDRESS1'), os.getenv('PRIVATE_KEY1')), Operations(os.getenv('ADDRESS2'), os.getenv('PRIVATE_KEY2')) ] # 为每个实例创建独立进程 processes = [multiprocessing.Process(target=instance.run) for instance in instances] for p in processes: p.start() for p in processes: p.join()
4. asyncio最佳实践
- 绝对避免在异步函数中调用同步阻塞代码,所有耗时操作必须异步化或放到线程/进程中;
- 用
asyncio.create_task正确创建任务,确保多个任务能被事件循环调度; - 不要滥用
await asyncio.sleep(0)强制切换任务,优先依赖真正的IO等待触发调度。
内容的提问来源于stack exchange,提问作者Lawnmoers The Epic Gamer

