詳解Python中asyncio.Queue長連接服務(wù)的內(nèi)存泄漏問題解決方法
問題背景
在基于 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.deque 有 maxlen 可以自動丟棄舊數(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_access 用 monotonic | 不受系統(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_TIMEOUT | 300s(5 min) | 取決于事件頻率。高頻事件可縮短到 60s,低頻可延長到 15min |
_CLEANUP_INTERVAL | 60s | 掃描開銷極?。ㄖ槐葧r(shí)間戳),可以更頻繁如 30s。但沒必要低于 10s |
對比其他語言的類似方案
| 語言/框架 | 類似機(jī)制 |
|---|---|
| Go | context.WithTimeout + select,超時(shí)自動退出 goroutine |
| Java | WeakReference + ReferenceQueue,GC 時(shí)回調(diào)清理 |
| Rust | Arc<Weak> + 定期 upgrade() 檢查 |
| Kubernetes | Lease 對象 + 定期續(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í)例代碼
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)程的教程,在多進(jìn)程編程中十分有用,示例代碼基于Python2.x版本,需要的朋友可以參考下2015-04-04
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)名稱方式,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧2020-01-01
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ù)器日志示例,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧2019-12-12
python 動態(tài)調(diào)用函數(shù)實(shí)例解析
這篇文章主要介紹了python 動態(tài)調(diào)用函數(shù)實(shí)例解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2019-10-10

