asyncio.create_task中误用await引发的意外并发及信号量失效问题解析
你好,我来帮你拆解这个问题里的两个核心疑惑:为什么明明在create_task里加了await,代码却还是全并发执行,同时设置的信号量完全没起到限制作用?咱们一步步来分析:
一、先搞懂task = asyncio.create_task(await self.fake_request(url))到底做了什么
你说得没错,这里把await放在create_task里面是完全错误的,但咱们得先理清这行代码的执行顺序:
- 首先执行
await self.fake_request(url):这会直接运行fake_request协程,直到它返回结果。 fake_request里的逻辑是:进入async with self.semaphore获取信号量,然后直接返回another_async_function(url)的协程对象(注意这里没有await),之后async with上下文退出,信号量立刻被释放。- 最后
create_task拿到的是another_async_function的协程对象,把它包装成任务丢进事件循环。
关键就在这里:await self.fake_request(url)的执行几乎是瞬间完成的——它只是短暂占用信号量,然后马上返回一个未执行的协程对象,信号量随即释放。所以next方法里的for循环不会被阻塞,能快速遍历所有7个URL,把所有another_async_function的任务都塞进事件循环,最终所有任务一起并发执行,这就是你看到所有Start task同时打印的原因。
二、信号量完全失效的根源
你设置的max_concurrent_requests=2本来是想限制同时执行的任务数,但信号量的作用范围完全错了:
- 信号量只在
fake_request执行的那几毫秒里被占用,而真正要限制的another_async_function(模拟网络请求的耗时操作)根本没被信号量包裹。 - 当
another_async_function的任务开始执行时,信号量早就被fake_request的async with上下文释放了,自然起不到任何并发限制作用。
三、代码里的两个核心错误点总结
错误1:
create_task中误用awaitasyncio.create_task的作用是把协程对象包装成可调度的任务,让事件循环异步执行它,不需要先await协程。你这里的写法相当于先同步执行完fake_request,再把返回的协程包装成任务,完全浪费了create_task的异步调度能力。如果fake_request是个耗时操作,你的for循环会被逐个阻塞,变成完全串行执行。错误2:
fake_request没有正确等待异步任务完成fake_request里应该await another_async_function(url),而不是直接返回它的协程对象。只有这样,async with semaphore的上下文才会一直持有信号量,直到another_async_function执行完毕,真正实现并发限制。
四、修正后的正确写法示例
import asyncio class RequestHandler: def __init__(self, max_concurrent_requests): self.semaphore = asyncio.Semaphore(max_concurrent_requests) async def another_async_function(self, url): print(f"Start task for {url}") await asyncio.sleep(5) print(f"End task for {url}") return url async def fake_request(self, url): async with self.semaphore: # 这里要await异步任务,让信号量在任务执行期间保持占用 return await self.another_async_function(url) async def next(self, urls): tasks = [] for url in urls: print(f"Scheduling task for {url}...") # create_task直接接收协程对象,不需要await task = asyncio.create_task(self.fake_request(url)) tasks.append(task) for task in tasks: result = await task yield result async def main(self): urls = ["url1", "url2", "url3", "url4", "url5", "url6", "url7"] async for result in self.next(urls): print(f"Processed result: {result}") request_handler = RequestHandler(max_concurrent_requests=2) asyncio.run(request_handler.main())
修正后你会看到,每次只有2个任务同时执行,完成后才会启动下一批,信号量的限制作用就生效了。
备注:内容来源于stack exchange,提问作者himanshu sharma

