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

asyncio.create_task中误用await引发的意外并发及信号量失效问题解析

asyncio.create_task中误用await引发的意外并发及信号量失效问题解析

你好,我来帮你拆解这个问题里的两个核心疑惑:为什么明明在create_task里加了await,代码却还是全并发执行,同时设置的信号量完全没起到限制作用?咱们一步步来分析:

一、先搞懂task = asyncio.create_task(await self.fake_request(url))到底做了什么

你说得没错,这里把await放在create_task里面是完全错误的,但咱们得先理清这行代码的执行顺序:

  1. 首先执行await self.fake_request(url):这会直接运行fake_request协程,直到它返回结果。
  2. fake_request里的逻辑是:进入async with self.semaphore获取信号量,然后直接返回another_async_function(url)的协程对象(注意这里没有await),之后async with上下文退出,信号量立刻被释放。
  3. 最后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中误用await
    asyncio.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 08:57:59