Python asyncio使用Semaphore限制并发执行数不生效问题求解
Python asyncio Semaphore 并发限制失效问题
问题代码
最初实现的限制并发的异步程序代码如下:
import asyncio async def gather_with_concurrency(*tasks,limit=3): semaphore = asyncio.Semaphore(limit) async def sem_task(task): async with semaphore: return await task return await asyncio.gather(*(sem_task(task) for task in tasks)) async def test(second): print(f"{second} start") await asyncio.sleep(second) print(f"{second} done") return second async def _run(): tasks = [] for i in range(5): tasks.append(asyncio.create_task(test(i+1))) return await gather_with_concurrency(*tasks) def mainloop(): loop = asyncio.get_event_loop() results = loop.run_until_complete(_run()) if __name__ == '__main__': mainloop()
预期输出为最多3个任务同时执行:
1 start 2 start 3 start 1 end 4 start 2 end 5 start 3 end 4 end 5 end
实际输出为所有任务同时启动,Semaphore限制完全失效:
1 start 2 start 3 start 4 start 5 start 1 done 2 done 3 done 4 done 5 done
问题排查
在sem_task的Semaphore上下文内新增打印后,输出如下:
1 start 2 start 3 start 4 start 5 start Starting <Task pending name='Task-2' coro=<test() running at stackoverflow.py:16> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7fe87b4ab460>()]>> Starting <Task pending name='Task-3' coro=<test() running at stackoverflow.py:16> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7fe87b4ab490>()]>> Starting <Task pending name='Task-4' coro=<test() running at stackoverflow.py:16> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7fe87b4ab4c0>()]>> 1 done Starting <Task pending name='Task-5' coro=<test() running at stackoverflow.py:16> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7fe87b4ab4f0>()]>> 2 done Starting <Task pending name='Task-6' coro=<test() running at stackoverflow.py:16> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7fe87b4ab520>()]>> 3 done 4 done 5 done
可以看到Semaphore确实限制了Starting打印的并发,但test函数的开头打印已经全部提前触发。
失效原因
核心问题出在任务创建时机:你在调用gather_with_concurrency之前,就已经通过asyncio.create_task(test(i+1))创建了所有任务。create_task被调用时,对应协程会立即被加入事件循环调度,不会等到sem_task里的await task才执行,所以Semaphore的限制逻辑还没生效,所有任务第一个await之前的同步代码(也就是print(f"{second} start"))就已经全部执行完毕了。
修复方案
不要提前创建任务,直接传入未执行的协程对象,在Semaphore的上下文内再创建任务/执行协程,确保限制逻辑在任务启动前生效,修复后代码如下:
import asyncio async def gather_with_concurrency(*cors,limit=3): semaphore = asyncio.Semaphore(limit) async def sem_task(cor): async with semaphore: print("Starting", cor) return await asyncio.create_task(cor) return await asyncio.gather(*(sem_task(cor) for cor in cors)) async def test(second): print(f"{second} start") await asyncio.sleep(second) print(f"{second} done") return second async def _run(): cors = [] for i in range(5): cors.append(test(i+1)) return await gather_with_concurrency(*cors) def mainloop(): loop = asyncio.get_event_loop() results = loop.run_until_complete(_run()) if __name__ == '__main__': mainloop()
关于create_task的疑问解答
调用asyncio.create_task()时,会立即将协程提交到当前事件循环的调度队列,只要事件循环获得执行权,就会优先执行协程中第一个await之前的所有同步代码,直到协程执行到第一个await语句主动让出执行权,事件循环才会切换去执行其他任务。
内容的提问来源于stack exchange,提问作者user20533
相关产品推荐
相关产品推荐

