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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 06:36:03