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

Python 3.12中队列使用后内存无法释放的问题求助

Python 3.12中队列使用后内存无法释放的问题求助

大伙好,我最近遇到一个头疼的内存问题,想请教下各位。

我原本想实现一个简单的流程:用协程通过HTTP流式拉取数据到队列,再开另一个协程从队列读数据写入磁盘。但实际运行时发现,队列的内存占用峰值一直降不下来——哪怕队列清空、甚至删除之后,进程还是死死占着大量内存,我的进程是长期运行的,这么耗内存肯定不行。

我的运行环境是Python 3.12.3,Linux内核6.8.0-51。下面是能复现问题的最小示例(原本我用的是multiprocessing队列,但换成asyncio队列也能复现同样的内存问题):

from __future__ import annotations

import asyncio
from pathlib import Path
from socket import AF_INET
import psutil

import aiofiles.os

from aiohttp import (
    ClientError,
    ClientSession,
    ClientTimeout,
    TCPConnector,
)


async def write_to_disk(dest: Path, queue: asyncio.Queue) -> None:
    async with aiofiles.open(dest, "wb") as f:
        while data := await queue.get():
            # 让出控制权给HTTP读取协程,模拟耗时操作
            await asyncio.sleep(0)
            await f.write(data)
            # 这个调用好像没作用
            queue.task_done()

        # 最后一个空数据是结束标记,这里的task_done也没效果
        queue.task_done()

    print("write to disk finished")


async def log_mem():
    while True:
        process = psutil.Process()
        used = round(process.memory_info().rss / 1024**2, 2)
        print(f"mem used: {used} MiB")
        await asyncio.sleep(5)


async def start_http_stream(
    queue: asyncio.Queue,
    connect_timeout: int = 3,
    read_timeout: int = 10,
) -> None:
    timeout = ClientTimeout(connect=connect_timeout, sock_read=read_timeout)
    conn = TCPConnector(family=AF_INET)

    url = "http://somedownloadurl.com/filetodownload"
    extra_params = {"json": {}}

    async with ClientSession(
        connector=conn,
    ) as session:
        try:
            async with session.post(
                url, timeout=timeout, **extra_params
            ) as resp:
                bytes_downloaded = 0
                percent_logged = 0
                content_length = int(resp.headers.get("Approx-Content-Length"))

                print(
                    "Download size:",
                    round(content_length / 1024 / 1024 / 1024, 2),
                    "GiB",
                )

                async for chunk, _ in resp.content.iter_chunks():
                    if not chunk:
                        # 会收到空chunk,如果传过去的话另一端会认为流结束
                        continue

                    bytes_downloaded += len(chunk)
                    percent = int((bytes_downloaded / content_length) * 100)
                    if percent > percent_logged and percent % 5 == 0:
                        percent_logged = percent
                        print(
                            f"Downloaded: {percent}%, q size: {queue.qsize()}"
                        )

                    # 用这个简单的方式做背压,之前用有界队列的时候遇到过读取超时的问题
                    if queue.qsize() >= 8192:
                        await asyncio.sleep(0.05)

                    await queue.put(chunk)
                # 结束标记
                await queue.put(b"")
        except (ClientError, OSError) as e:
            print("EXC", e)

    print("Http stream finished")


async def main():
    q = asyncio.Queue()

    stream_coro = start_http_stream(q)
    write_coro = write_to_disk("file_download.tar", q)
    log_mem_task = asyncio.create_task(log_mem())

    await asyncio.gather(stream_coro, write_coro)
    print("tasks finished... sleeping 30 seconds")
    await asyncio.sleep(30)
    print("Current q size:", q.qsize())
    print("deleting queue and sleeping 30 seconds")
    del q
    await asyncio.sleep(30)
    log_mem_task.cancel()
    try:
        await log_mem_task
    except asyncio.CancelledError:
        print("log mem cancelled")


asyncio.run(main())

下面是运行输出(中间部分我精简掉了):

mem used: 31.62 MiB
Download size: 36.09 GiB
Downloaded: 5%, q size: 5799
mem used: 1534.69 MiB
mem used: 2049.82 MiB
Downloaded: 10%, q size: 8168
mem used: 2049.82 MiB
Downloaded: 15%, q size: 8175
mem used: 2048.55 MiB
Downloaded: 20%, q size: 8180
<snip>
mem used: 2052.88 MiB
mem used: 2052.88 MiB
Downloaded: 80%, q size: 8113
mem used: 2052.88 MiB
Downloaded: 85%, q size: 8112
mem used: 2052.88 MiB
Downloaded: 90%, q size: 8148
mem used: 2052.45 MiB
mem used: 2052.45 MiB
Downloaded: 95%, q size: 8042
mem used: 2052.45 MiB
Downloaded: 100%, q size: 8179
Http stream finished
mem used: 2052.45 MiB
write to disk finished
tasks finished... sleeping 30 seconds
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
Current q size: 0
deleting queue and sleeping 30 seconds
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
mem used: 2049.34 MiB
log mem cancelled

从输出能看到,哪怕所有任务都完成了,队列大小变成0,甚至删除队列之后,内存还是没释放。我用tracemalloc查过,堆内存其实已经释放了,这是不是CPython的某种机制导致的?

我实在不知道该怎么进一步调试或者释放内存了,都快被逼得要在下载完成后重启应用了,这显然不是个好方案。有没有大佬能指点一下我该怎么做?


补充说明

  1. 我试过手动调用gc.collect(),结果还是一样:
Taking out the trash and sleeping 30 seconds
mem used: 1990.68 MiB
mem used: 1990.68 MiB
mem used: 1990.68 MiB
mem used: 1990.68 MiB
mem used: 1990.68 MiB
mem used: 1990.68 MiB
log mem cancelled
  1. 这里指的是实际物理内存,不是虚拟内存。
  2. 队列删除后的内存占用情况:
    队列删除后的内存占用
  3. 进程退出后的内存释放情况:
    进程退出后的内存释放

备注:内容来源于stack exchange,提问作者David White

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 12:23:00