最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

詳解Python中asyncio.Queue長連接服務(wù)的內(nèi)存泄漏問題解決方法

 更新時(shí)間:2026年05月28日 08:54:52   作者:十有八七  
這篇文章主要為大家詳細(xì)介紹了Python中asyncio.Queue長連接服務(wù)的內(nèi)存泄漏問題解決方法,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以了解一下

問題背景

在基于 asyncio 的長連接服務(wù)(如 SSE 推送、WebSocket、實(shí)時(shí)儀表盤)中,事件總線是最常見的基礎(chǔ)設(shè)施之一。核心模式很簡單:

# 訂閱者注冊一個(gè) asyncio.Queue
queue = event_bus.subscribe("task:progress")

# 發(fā)布者往所有匹配的 Queue 里 put_nowait
event_bus.emit(Event("task:progress", data={"percent": 80}))

但生產(chǎn)環(huán)境中,訂閱者的生命周期往往不可控——前端斷連、瀏覽器關(guān)閉、SSE 連接超時(shí)、客戶端忘記調(diào)用 unsubscribe……這些場景都會導(dǎo)致一個(gè)結(jié)果:

Queue 永遠(yuǎn)不會被消費(fèi),也永遠(yuǎn)不會被移除。

每次新連接創(chuàng)建一個(gè) Queue,斷連后 Queue 就成了"孤兒"——仍然掛在訂閱列表里,emit() 時(shí)仍會往里面塞數(shù)據(jù)。日積月累,內(nèi)存持續(xù)增長,最終 OOM。

這就是 長連接服務(wù)的隱形內(nèi)存泄漏:不是你的代碼寫錯(cuò)了,而是你沒有為"不正常的離開"兜底。

錯(cuò)誤方案與它們的缺陷

方案 A:依賴客戶端主動 unsubscribe

# SSE 端點(diǎn)
async def sse_endpoint(request):
    queue = event_bus.subscribe("task:progress")
    try:
        while True:
            event = await queue.get()
            yield sse_format(event)
    finally:
        event_bus.unsubscribe(queue)

問題finally 只在協(xié)程正常退出時(shí)執(zhí)行。如果客戶端直接斷開 TCP 連接,服務(wù)端的 queue.get() 可能還在 await 中,協(xié)程不會立刻被取消。即使框架做了清理,也存在時(shí)間窗口——連接斷開到協(xié)程被取消之間的幾百毫秒內(nèi),emit() 已經(jīng)往死 Queue 里塞了若干事件。

更關(guān)鍵的是:你無法假設(shè)所有調(diào)用方都記得 unsubscribe。EventBus 是基礎(chǔ)設(shè)施代碼,它必須對自己的健康負(fù)責(zé)。

方案 B:給 Queue 設(shè) maxlen

queue = asyncio.Queue(maxsize=100)

問題maxsize 限制的是 put 的阻塞,不是內(nèi)存。用 put_nowait 時(shí),超出會直接拋 QueueFull 異常——你要么吞掉異常(丟事件),要么讓 emit 失?。ㄓ绊懰杏嗛喺撸6?maxlen=0(無限)是默認(rèn)行為,改不了。

collections.dequemaxlen 可以自動丟棄舊數(shù)據(jù),但 asyncio.Queue 沒有。

方案 C:弱引用 WeakRef

self._subscribers.append(weakref.ref(queue))

問題asyncio.Queue 一旦沒有強(qiáng)引用就會立刻被 GC 回收——但 Queue 本身就是被 EventBus 持有的,你恰恰需要 EventBus 來引用它。弱引用在這里邏輯上自相矛盾。

正確方案:TrackedQueue + 后臺清理協(xié)程

核心思路:讓 Queue 自己告訴系統(tǒng)"我還活著",系統(tǒng)定期檢查誰已經(jīng)沉默太久

第一步:包裝 Queue,記錄最后活躍時(shí)間

class _TrackedQueue(asyncio.Queue[Event]):
    """帶追蹤信息的 Queue,記錄最后消費(fèi)時(shí)間和訂閱參數(shù)"""

    __slots__ = ("last_access", "_sub_event_type", "_sub_task_id")

    def __init__(
        self,
        event_type: str | None = None,
        task_id: uuid.UUID | None = None,
    ) -> None:
        super().__init__()
        self.last_access: float = time.monotonic()
        self._sub_event_type = event_type
        self._sub_task_id = task_id

    def touch(self) -> None:
        self.last_access = time.monotonic()

關(guān)鍵設(shè)計(jì)

設(shè)計(jì)點(diǎn)說明
繼承 asyncio.Queue對消費(fèi)方完全透明,await queue.get() 無需任何改動
__slots__避免每個(gè)實(shí)例創(chuàng)建 __dict__,大量 Queue 場景下節(jié)省內(nèi)存
last_accessmonotonic不受系統(tǒng)時(shí)間回調(diào)影響,比 time.time() 更安全
保存訂閱參數(shù)清理時(shí)需要知道"這個(gè) Queue 掛在哪個(gè)列表里",否則要從所有列表里線性搜索

第二步:全局注冊表 + 統(tǒng)一清理入口

class EventBus:
    def __init__(self) -> None:
        self._subscribers: dict[str, list[_TrackedQueue]] = defaultdict(list)
        self._global_subscribers: list[_TrackedQueue] = []
        self._task_subscribers: dict[uuid.UUID, list[_TrackedQueue]] = defaultdict(list)
        self._all_queues: set[_TrackedQueue] = set()  # ← 全局注冊表
        self._cleanup_task: asyncio.Task | None = None

    def _remove_queue(self, q: _TrackedQueue) -> None:
        """從所有訂閱列表中移除隊(duì)列"""
        self._all_queues.discard(q)
        if q._sub_task_id is not None:
            subs = self._task_subscribers.get(q._sub_task_id, [])
            if q in subs:
                subs.remove(q)
                if not subs:
                    del self._task_subscribers[q._sub_task_id]
        elif q._sub_event_type is None:
            if q in self._global_subscribers:
                self._global_subscribers.remove(q)
        else:
            subs = self._subscribers.get(q._sub_event_type, [])
            if q in subs:
                subs.remove(q)

為什么需要 _all_queues

訂閱者按 event_type / task_id / 全局分三組存儲。如果沒有全局注冊表,清理時(shí)需要遍歷所有分組查找——O(n) 且代碼丑陋。_all_queues 是一個(gè) set,查找 O(1),遍歷一遍就能找出所有閑置 Queue。

_remove_queue 根據(jù) Queue 自身保存的訂閱參數(shù),精準(zhǔn)定位它所在的分組列表,一次調(diào)用搞定所有清理。

第三步:后臺清理協(xié)程(懶啟動)

_QUEUE_IDLE_TIMEOUT = 300  # 5 分鐘
_CLEANUP_INTERVAL = 60      # 60 秒檢查一次

def _start_cleanup_loop(self) -> None:
    """懶啟動后臺清理協(xié)程"""
    if self._cleanup_task is not None:
        return
    try:
        loop = asyncio.get_running_loop()
    except RuntimeError:
        return
    self._cleanup_task = loop.create_task(self._cleanup_loop())

async def _cleanup_loop(self) -> None:
    while True:
        await asyncio.sleep(_CLEANUP_INTERVAL)
        self._cleanup_idle_queues()

def _cleanup_idle_queues(self) -> None:
    now = time.monotonic()
    stale = [q for q in self._all_queues if now - q.last_access > _QUEUE_IDLE_TIMEOUT]
    for q in stale:
        self._remove_queue(q)

懶啟動:在 subscribe() 首次被調(diào)用時(shí)才啟動清理協(xié)程。避免在 import 時(shí)或無訂閱者的場景下創(chuàng)建無用的 Task。

第四步:在 emit 中 touch

def emit(self, event: Event) -> None:
    if event.task_id is not None:
        for queue in self._task_subscribers.get(event.task_id, []):
            queue.touch()  # ← 標(biāo)記活躍
            queue.put_nowait(event)

    for queue in self._subscribers.get(event.event_type, []):
        queue.touch()
        queue.put_nowait(event)

    for queue in self._global_subscribers:
        queue.touch()
        queue.put_nowait(event)

touch() 的時(shí)機(jī)選擇:在 emit 時(shí)調(diào)用,而非在 queue.get() 時(shí)。

為什么?因?yàn)?emit 是 EventBus 自己的代碼路徑,完全可控。而 queue.get() 發(fā)生在消費(fèi)方的協(xié)程里,EventBus 無法插手也不應(yīng)該侵入。emit 時(shí) touch 的語義是:"有人還在往這個(gè) Queue 塞數(shù)據(jù)"——如果連 emit 都不再光顧這個(gè) Queue 了,說明訂閱的事件類型/任務(wù)已經(jīng)沒有新事件產(chǎn)生,Queue 的閑置是合理的。

如果需要更嚴(yán)格的活躍判定(消費(fèi)者是否真的在讀?。梢栽?events() 異步迭代器的 yield 前也加 touch()。

完整數(shù)據(jù)流

subscribe()                    emit()                   cleanup_loop()
    │                             │                          │
    ├─ _TrackedQueue()            ├─ queue.touch()           ├─ sleep(60s)
    ├─ _all_queues.add(q)         ├─ queue.put_nowait(e)     ├─ scan _all_queues
    ├─ 注冊到分組列表              │                          ├─ 找 idle > 300s
    └─ _start_cleanup_loop()      │                          └─ _remove_queue(q)

可觀測性:stats 屬性

@property
def stats(self) -> dict[str, int]:
    return {
        "global_subscribers": len(self._global_subscribers),
        "type_subscribers": sum(len(v) for v in self._subscribers.values()),
        "task_subscribers": sum(len(v) for v in self._task_subscribers.values()),
        "total_queues": len(self._all_queues),
    }

暴露給 /api/stats 或 Prometheus,可以直接觀察到:

  • total_queues 持續(xù)增長 → 清理協(xié)程可能沒啟動
  • total_queues 突降 → 可能是 idle_timeout 設(shè)得太短
  • 分組不均衡 → 某個(gè) task_id 或 event_type 訂閱者過多

模式總結(jié)

元素作用對應(yīng)概念
_TrackedQueue包裝原始 Queue,附加追蹤元數(shù)據(jù)裝飾器模式(繼承式)
_all_queues set全局注冊表,O(1) 查找 + 統(tǒng)一遍歷索引/注冊表
touch()標(biāo)記"我還在用"心跳/?;?/strong>
_cleanup_loop后臺定期掃描并回收GC 守護(hù)線程(協(xié)程版)
懶啟動首次 subscribe 時(shí)才創(chuàng)建 Task延遲初始化
訂閱參數(shù)保存在 Queue清理時(shí)精準(zhǔn)定位列表自描述對象

這個(gè)模式的通用性

不限于 EventBus。任何"生產(chǎn)者 → Queue → 消費(fèi)者"模型中,如果消費(fèi)者可能靜默退出(網(wǎng)絡(luò)斷連、進(jìn)程崩潰、邏輯錯(cuò)誤),都可以用這個(gè)模式:

  • 消息隊(duì)列客戶端:本地 Queue 緩存消費(fèi)失敗的場景
  • 實(shí)時(shí)推送網(wǎng)關(guān):SSE / WebSocket 的 per-connection Queue
  • 日志聚合:每個(gè)日志流一個(gè) Queue,流不再活躍時(shí)自動回收
  • 任務(wù)調(diào)度器:worker 的任務(wù)隊(duì)列,worker 掉線后回收未執(zhí)行任務(wù)

關(guān)鍵參數(shù)調(diào)優(yōu)

參數(shù)建議值調(diào)優(yōu)思路
_QUEUE_IDLE_TIMEOUT300s(5 min)取決于事件頻率。高頻事件可縮短到 60s,低頻可延長到 15min
_CLEANUP_INTERVAL60s掃描開銷極?。ㄖ槐葧r(shí)間戳),可以更頻繁如 30s。但沒必要低于 10s

對比其他語言的類似方案

語言/框架類似機(jī)制
Gocontext.WithTimeout + select,超時(shí)自動退出 goroutine
JavaWeakReference + ReferenceQueue,GC 時(shí)回調(diào)清理
RustArc<Weak> + 定期 upgrade() 檢查
KubernetesLease 對象 + 定期續(xù)約,超時(shí)自動刪除

Python 沒有 Go 的 context 取消鏈,沒有 Java 的 GC 回調(diào),沒有 Rust 的所有權(quán)。Python 的方式是顯式的:你自己寫清理邏輯,自己保證執(zhí)行。 這也是這個(gè)模式存在的意義。

附錄:完整實(shí)現(xiàn)

參考源碼:

"""事件總線 — asyncio.Queue 實(shí)現(xiàn)的發(fā)布/訂閱

支持:
- 發(fā)布事件(emit)
- 訂閱事件(subscribe)
- 按 event_type 過濾
- 按 task_id 過濾
- 超時(shí)隊(duì)列自動清理(防止內(nèi)存泄漏)
"""

from __future__ import annotations

import asyncio
import time
import uuid
from collections import defaultdict
from dataclasses import dataclass, field
from typing import Any, AsyncIterator

# 隊(duì)列最大閑置時(shí)間(秒),超過此時(shí)間未被消費(fèi)則視為泄漏
_QUEUE_IDLE_TIMEOUT = 300  # 5 minutes


@dataclass
class Event:
    """通用事件"""
    event_type: str
    data: dict[str, Any] = field(default_factory=dict)
    source: str | None = None  # 事件來源(如 node_id)
    task_id: uuid.UUID | None = None  # 所屬任務(wù)


class _TrackedQueue(asyncio.Queue[Event]):
    """帶追蹤信息的 Queue,記錄最后消費(fèi)時(shí)間和訂閱參數(shù)"""

    __slots__ = ("last_access", "_sub_event_type", "_sub_task_id")

    def __init__(
        self,
        event_type: str | None = None,
        task_id: uuid.UUID | None = None,
    ) -> None:
        super().__init__()
        self.last_access: float = time.monotonic()
        self._sub_event_type = event_type
        self._sub_task_id = task_id

    def touch(self) -> None:
        self.last_access = time.monotonic()


class EventBus:
    """異步事件總線"""

    def __init__(self) -> None:
        self._subscribers: dict[str, list[_TrackedQueue]] = defaultdict(list)
        self._global_subscribers: list[_TrackedQueue] = []
        self._task_subscribers: dict[uuid.UUID, list[_TrackedQueue]] = defaultdict(list)
        self._all_queues: set[_TrackedQueue] = set()
        self._cleanup_task: asyncio.Task | None = None  # type: ignore[type-arg]

    def _start_cleanup_loop(self) -> None:
        """啟動后臺清理協(xié)程(懶啟動)"""
        if self._cleanup_task is not None:
            return
        try:
            loop = asyncio.get_running_loop()
        except RuntimeError:
            return
        self._cleanup_task = loop.create_task(self._cleanup_loop())

    async def _cleanup_loop(self) -> None:
        """定期清理超時(shí)隊(duì)列"""
        while True:
            await asyncio.sleep(60)  # 每 60 秒檢查一次
            self._cleanup_idle_queues()

    def _cleanup_idle_queues(self) -> None:
        """清理閑置超時(shí)的隊(duì)列"""
        now = time.monotonic()
        stale = [q for q in self._all_queues if now - q.last_access > _QUEUE_IDLE_TIMEOUT]
        for q in stale:
            self._remove_queue(q)

    def _remove_queue(self, q: _TrackedQueue) -> None:
        """從所有訂閱列表中移除隊(duì)列"""
        self._all_queues.discard(q)
        if q._sub_task_id is not None:
            subs = self._task_subscribers.get(q._sub_task_id, [])
            if q in subs:
                subs.remove(q)
                if not subs:
                    del self._task_subscribers[q._sub_task_id]
        elif q._sub_event_type is None:
            if q in self._global_subscribers:
                self._global_subscribers.remove(q)
        else:
            subs = self._subscribers.get(q._sub_event_type, [])
            if q in subs:
                subs.remove(q)

    def emit(self, event: Event) -> None:
        """發(fā)布事件到所有匹配的訂閱者"""
        # 按 task_id 訂閱(最高優(yōu)先級,精準(zhǔn)匹配)
        if event.task_id is not None:
            for queue in self._task_subscribers.get(event.task_id, []):
                queue.touch()
                queue.put_nowait(event)

        # 按類型訂閱
        for queue in self._subscribers.get(event.event_type, []):
            queue.touch()
            queue.put_nowait(event)
        # 全局訂閱
        for queue in self._global_subscribers:
            queue.touch()
            queue.put_nowait(event)

    def subscribe(
        self,
        event_type: str | None = None,
        task_id: uuid.UUID | None = None,
    ) -> asyncio.Queue[Event]:
        """訂閱事件

        Args:
            event_type: 事件類型,None 表示不按類型過濾
            task_id: 任務(wù)ID,指定后只接收該任務(wù)的事件

        Returns:
            asyncio.Queue,消費(fèi)方通過 await queue.get() 獲取事件
        """
        self._start_cleanup_loop()
        queue = _TrackedQueue(event_type=event_type, task_id=task_id)
        self._all_queues.add(queue)
        if task_id is not None:
            self._task_subscribers[task_id].append(queue)
        elif event_type is None:
            self._global_subscribers.append(queue)
        else:
            self._subscribers[event_type].append(queue)
        return queue

    def unsubscribe(
        self,
        queue: asyncio.Queue[Event],
        event_type: str | None = None,
        task_id: uuid.UUID | None = None,
    ) -> None:
        """取消訂閱"""
        if isinstance(queue, _TrackedQueue):
            self._remove_queue(queue)
            return
        # fallback: 對非 TrackedQueue 的老式調(diào)用
        if task_id is not None:
            subs = self._task_subscribers.get(task_id, [])
            if queue in subs:
                subs.remove(queue)
                if not subs:
                    del self._task_subscribers[task_id]
        elif event_type is None:
            if queue in self._global_subscribers:
                self._global_subscribers.remove(queue)
        else:
            subs = self._subscribers.get(event_type, [])
            if queue in subs:
                subs.remove(queue)

    async def events(
        self,
        event_type: str | None = None,
        task_id: uuid.UUID | None = None,
    ) -> AsyncIterator[Event]:
        """異步迭代訂閱的事件"""
        queue = self.subscribe(event_type, task_id)
        try:
            while True:
                event = await queue.get()
                if isinstance(queue, _TrackedQueue):
                    queue.touch()
                yield event
        finally:
            self.unsubscribe(queue, event_type, task_id)

    def clear(self) -> None:
        """清空所有訂閱"""
        self._subscribers.clear()
        self._global_subscribers.clear()
        self._task_subscribers.clear()
        self._all_queues.clear()

    @property
    def stats(self) -> dict[str, int]:
        """返回當(dāng)前訂閱統(tǒng)計(jì)"""
        return {
            "global_subscribers": len(self._global_subscribers),
            "type_subscribers": sum(len(v) for v in self._subscribers.values()),
            "task_subscribers": sum(len(v) for v in self._task_subscribers.values()),
            "total_queues": len(self._all_queues),
        }


# 全局事件總線實(shí)例
_event_bus: EventBus | None = None


def get_event_bus() -> EventBus:
    """獲取全局事件總線"""
    global _event_bus
    if _event_bus is None:
        _event_bus = EventBus()
    return _event_bus

以上就是詳解Python中asyncio.Queue長連接服務(wù)的內(nèi)存泄漏問題解決方法的詳細(xì)內(nèi)容,更多關(guān)于Python asyncio.Queue內(nèi)存泄漏問題解決的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Python繪制loss曲線和準(zhǔn)確率曲線實(shí)例代碼

    Python繪制loss曲線和準(zhǔn)確率曲線實(shí)例代碼

    pytorch雖然使用起來很方便,但在一點(diǎn)上并沒有tensorflow方便,就是繪制模型訓(xùn)練時(shí)在訓(xùn)練集和驗(yàn)證集上的loss和accuracy曲線(共四條),下面這篇文章主要給大家介紹了關(guān)于Python繪制loss曲線和準(zhǔn)確率曲線的相關(guān)資料,需要的朋友可以參考下
    2022-08-08
  • 在Python程序中實(shí)現(xiàn)分布式進(jìn)程的教程

    在Python程序中實(shí)現(xiàn)分布式進(jìn)程的教程

    這篇文章主要介紹了在Python程序中實(shí)現(xiàn)分布式進(jìn)程的教程,在多進(jìn)程編程中十分有用,示例代碼基于Python2.x版本,需要的朋友可以參考下
    2015-04-04
  • python中子類調(diào)用父類函數(shù)的方法示例

    python中子類調(diào)用父類函數(shù)的方法示例

    Python中類的初始化方法是__init__(),因此父類、子類的初始化方法都是這個(gè),下面這篇文章主要給大家介紹了關(guān)于python中子類調(diào)用父類函數(shù)的方法示例,文中通過示例代碼介紹的非常詳細(xì),需要的朋友可以參考下。
    2017-08-08
  • TensorFlow查看輸入節(jié)點(diǎn)和輸出節(jié)點(diǎn)名稱方式

    TensorFlow查看輸入節(jié)點(diǎn)和輸出節(jié)點(diǎn)名稱方式

    今天小編就為大家分享一篇TensorFlow查看輸入節(jié)點(diǎn)和輸出節(jié)點(diǎn)名稱方式,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-01-01
  • Python?Base64編碼和解碼操作

    Python?Base64編碼和解碼操作

    Base64?就是一種基于64個(gè)可打印字符來表示二進(jìn)制數(shù)據(jù)的方法,這篇文章主要介紹了Python?Base64編碼和解碼,需要的朋友可以參考下
    2022-12-12
  • python操作redis方法總結(jié)

    python操作redis方法總結(jié)

    本篇文章給大家總結(jié)了python操作redis的實(shí)際方法和實(shí)例代碼,有興趣的朋友參考學(xué)習(xí)下。
    2018-06-06
  • Python輸出由1,2,3,4組成的互不相同且無重復(fù)的三位數(shù)

    Python輸出由1,2,3,4組成的互不相同且無重復(fù)的三位數(shù)

    這篇文章主要介紹了Python輸出由1,2,3,4組成的互不相同且無重復(fù)的三位數(shù),分享了相關(guān)代碼示例,小編覺得還是挺不錯(cuò)的,具有一定借鑒價(jià)值,需要的朋友可以參考下
    2018-02-02
  • python實(shí)現(xiàn)tail實(shí)時(shí)查看服務(wù)器日志示例

    python實(shí)現(xiàn)tail實(shí)時(shí)查看服務(wù)器日志示例

    今天小編就為大家分享一篇python實(shí)現(xiàn)tail實(shí)時(shí)查看服務(wù)器日志示例,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2019-12-12
  • python 動態(tài)調(diào)用函數(shù)實(shí)例解析

    python 動態(tài)調(diào)用函數(shù)實(shí)例解析

    這篇文章主要介紹了python 動態(tài)調(diào)用函數(shù)實(shí)例解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-10-10
  • python中jieba模塊的深入了解

    python中jieba模塊的深入了解

    這篇文章主要介紹了python中jieba模塊的深入了解,jieba模塊是一個(gè)python第三方中文分詞模塊,可以用于將語句中的中文詞語分離出來
    2022-06-06

最新評論

会昌县| 苗栗市| 鹿邑县| 蓝山县| 乌兰察布市| 泰顺县| 基隆市| 双辽市| 临汾市| 宁夏| 上高县| 保定市| 长汀县| 文化| 类乌齐县| 伊金霍洛旗| 漯河市| 汪清县| 广安市| 云阳县| 桦甸市| 香格里拉县| 吴桥县| 台州市| 庄浪县| 建宁县| 将乐县| 兴安县| 清水县| 句容市| 巫溪县| 科尔| 高安市| 五河县| 湖北省| 东安县| 小金县| 潼关县| 通道| 镇原县| 申扎县|