Flask结合SocketIO时MongoDB Change Stream的正确关闭方法咨询
解决方案
核心思路
利用弱引用字典存储Change Stream实例,避免强引用导致的内存泄漏,同时将流与Socket连接绑定,连接断开时自动触发清理,既不需要手动维护全局映射表,也能保证所有流正确关闭。
MongoDB的Change Stream本身提供close()方法,调用后迭代会立刻终止,只需要在正确的时机触发该方法即可。
代码实现
首先引入弱引用字典存储流实例,弱引用不会阻止GC回收对象,连接释放后对应流会自动从字典中移除,不会产生内存残留:
import weakref from flask import request # 弱引用字典,仅持有流实例的弱引用,无其他强引用时自动回收 user_streams = weakref.WeakValueDictionary()
修改监听逻辑,增加异常捕获与资源清理逻辑:
@socket.on('listen_user') def listen(uid): print(f"Listening to user : {uid}", flush=True) stream = db.listen_user(uid) # 用uid+当前连接的sid作为key,支持同一用户多端同时监听 stream_key = f"{uid}:{request.sid}" user_streams[stream_key] = stream try: for change in stream: socket.emit("listen_user_callback", change) except Exception: # 流被主动关闭时会抛出异常,直接退出即可 pass finally: # 兜底保证流被关闭 if not stream.closed: stream.close() user_streams.pop(stream_key, None)
新增Socket断开事件处理,客户端主动断开或网络异常时自动关闭对应流:
@socket.on('disconnect') def handle_disconnect(): sid = request.sid # 找到当前连接关联的所有流并关闭 related_keys = [k for k in user_streams.keys() if k.endswith(f":{sid}")] for key in related_keys: stream = user_streams.get(key) if stream and not stream.closed: stream.close()
完成TODO中的停止监听接口逻辑:
@app.route('/users/<uid>/stop', methods=["GET"]) def stop_listener(uid): print(f"Stop listening to user : {uid}", flush=True) # 关闭该用户所有打开的监听流 related_keys = [k for k in user_streams.keys() if k.startswith(f"{uid}:")] for key in related_keys: stream = user_streams.get(key) if stream and not stream.closed: stream.close() return {"status": "success"}
注意事项
- 如果需要支持关闭指定设备的监听流,可以在调用停止接口时传入对应连接的sid,匹配key时同时校验uid和sid即可。
- 多实例部署场景下,HTTP停止请求可能会打到其他实例,可通过Redis发布订阅广播关闭信号,所有实例收到信号后关闭自身存储的对应用户流即可。
- 弱引用字典的存储完全不会导致内存泄漏,只有处于活跃监听状态的流会被暂存,流关闭后会自动被GC回收,无需手动清理过期条目。
内容的提问来源于stack exchange,提问作者Tom3652
相关产品推荐
相关产品推荐

