Python使用asyncio Semaphore时内存持续增长,移除后正常该如何解决?
问题根因
- 你遇到的内存上涨问题本质是任务生产速度远大于消费速度,Semaphore仅能限制同时处于执行状态的任务数量,不会限制排队等待的任务数量。按你给出的示例参数计算:每0.01秒生成1个任务,单个任务执行耗时5秒,单账号2并发的情况下每秒最多只能完成0.4个任务,每秒会新增99.6个排队任务,这些等待执行的任务都是Python对象,会持续占用内存,最终导致内存持续升高。
- 移除Semaphore后任务不需要排队,创建后直接进入执行状态,任务执行完成后会被GC自动回收,整体任务存量会稳定在
(1/0.01)*5 = 500个左右,所以不会出现内存持续增长的问题。但如果后续任务执行耗时变长,就算没有Semaphore同样会出现任务堆积内存上涨的问题。
解决方案
不需要完全重构整体代码结构,只需要调整任务生成逻辑,添加排队长度限制即可,推荐用asyncio.Queue实现生产者消费者模式,从根源上避免任务无限堆积:
import asyncio async def task(acc_id): print(f"acc_id {acc_id} received the task") await asyncio.sleep(5) print(f"acc_id {acc_id} finished the task") # 消费者协程:固定并发数,持续从队列取任务执行 async def consumer(acc_id, queue): while True: task_item = await queue.get() try: await task_item finally: queue.task_done() async def create_tasks(acc_id): # 单账号最多允许10个等待执行的任务,可根据实际业务场景调整阈值 queue = asyncio.Queue(maxsize=10) # 启动2个消费者,对应原来的2并发限制 for _ in range(2): asyncio.create_task(consumer(acc_id, queue)) while True: # 队列满时put会自动阻塞,不会继续生成新任务 await queue.put(task(acc_id)) await asyncio.sleep(0.01) # request per 0.01 sec async def create_tasks_for_all_accs(): for acc_id in range(3): asyncio.create_task(create_tasks(acc_id)) def main(): event_loop = asyncio.get_event_loop() asyncio.ensure_future(create_tasks_for_all_accs()) event_loop.run_forever() if __name__ == "__main__": main()
代码调整说明
- 取消了原有的Semaphore逻辑,改用固定数量的消费者协程实现并发控制,和原有的并发限制逻辑效果完全一致
- 每个账号对应一个设置了最大长度的任务队列,队列满时生产者会自动阻塞,不会无限生成新任务,从根源上避免了任务堆积
- 移除了不必要的event_loop参数传递,3.7+版本的asyncio可以直接通过
asyncio.create_task创建任务,不需要手动传入loop - 不建议直接读取
BoundedSemaphore的_waiters私有属性判断排队长度,不同Python版本内部实现可能会变化,稳定性没有保障
内容的提问来源于stack exchange,提问作者zevcc
相关产品推荐
相关产品推荐

