Flask服务如何保持WebSocket长连接复用处理入站HTTP请求
实现方案
你可以通过单例模式全局管理WebSocket连接,搭配请求ID匹配、自动重连、心跳保活机制实现连接复用,具体实现逻辑如下:
核心改造点
- 避免每次请求新建连接,全局维护唯一一个WebSocket连接实例
- 给每个请求分配唯一请求ID,响应时通过ID匹配对应请求,解决全双工通信下响应乱序问题
- 增加心跳和自动重连逻辑,保证连接长期可用
- 适配Flask异步运行环境,用ASGI服务器(如Uvicorn)启动服务以支持异步视图
完整代码示例
import asyncio import json import uuid from flask import Flask, jsonify import websockets app = Flask(__name__) class WebSocketClient: _instance = None # 单例模式,全局唯一实例 def __new__(cls, *args, **kwargs): if not cls._instance: cls._instance = super().__new__(cls) return cls._instance def __init__(self, ws_uri): if hasattr(self, 'ws_uri'): return self.ws_uri = ws_uri self.websocket = None # 存储待响应的请求:key是请求ID,value是asyncio.Future对象 self.pending_requests = {} # 连接状态标记 self.connected = False # 保活任务、接收消息任务引用 self.keepalive_task = None self.recv_task = None async def connect(self): """建立WebSocket连接,断开时自动重连""" while True: try: self.websocket = await websockets.connect(self.ws_uri) self.connected = True print("WebSocket连接建立成功") # 启动接收消息循环和心跳任务 self.recv_task = asyncio.create_task(self._recv_loop()) self.keepalive_task = asyncio.create_task(self._keepalive()) # 等待连接断开 await self.websocket.wait_closed() except Exception as e: print(f"WebSocket连接断开,5秒后重试:{str(e)}") finally: self.connected = False # 清理现有任务 if self.recv_task: self.recv_task.cancel() if self.keepalive_task: self.keepalive_task.cancel() # 清空待处理请求,返回异常 for future in self.pending_requests.values(): if not future.done(): future.set_exception(RuntimeError("WebSocket连接断开")) self.pending_requests.clear() await asyncio.sleep(5) async def _recv_loop(self): """持续接收服务端返回的消息,匹配对应请求""" while self.connected: try: raw_data = await self.websocket.recv() data = json.loads(raw_data) req_id = data.get("req_id") if req_id and req_id in self.pending_requests: future = self.pending_requests.pop(req_id) if not future.done(): future.set_result(data) except Exception as e: print(f"接收消息出错:{str(e)}") break async def _keepalive(self): """心跳保活,每30秒发一次ping""" while self.connected: try: await self.websocket.ping() await asyncio.sleep(30) except Exception as e: print(f"心跳失败:{str(e)}") break async def send_request(self, payload): """发送请求并等待响应""" if not self.connected: raise RuntimeError("WebSocket连接未就绪") # 生成唯一请求ID req_id = str(uuid.uuid4()) payload["req_id"] = req_id # 创建Future对象等待结果 future = asyncio.get_event_loop().create_future() self.pending_requests[req_id] = future try: await self.websocket.send(json.dumps(payload).encode("utf-8")) # 等待结果,超时30秒 return await asyncio.wait_for(future, timeout=30) except Exception as e: if req_id in self.pending_requests: self.pending_requests.pop(req_id) raise e # 初始化全局WebSocket客户端,替换为你的第三方服务地址 ws_client = WebSocketClient("wss://the-third-party-server.com/xyz") # Flask启动时自动运行WebSocket连接任务 @app.before_first_request def start_ws_client(): loop = asyncio.get_event_loop() loop.create_task(ws_client.connect()) # 异步视图示例 @app.route("/api/request", methods=["POST"]) async def handle_request(): # 这里可以拿到HTTP请求的参数构造payload payload = {"test": "data"} try: result = await ws_client.send_request(payload) return jsonify(result) except Exception as e: return jsonify({"error": str(e)}), 500 if __name__ == "__main__": # 用uvicorn启动Flask ASGI应用 import uvicorn uvicorn.run(app, host="0.0.0.0", port=5000)
注意事项
- 如果你使用的是传统WSGI服务器,需要单独启动一个后台线程运行asyncio事件循环来管理WebSocket连接,避免阻塞Flask的同步请求处理
- 第三方服务的响应必须携带你发送的
req_id参数,否则无法匹配对应请求,如果第三方服务不支持自定义参数,你需要调整匹配逻辑(比如按请求顺序排队处理,同一时间只处理一个请求) - 可以根据实际需求调整心跳间隔、请求超时时间、重连等待时间参数
内容的提问来源于stack exchange,提问作者Joko
相关产品推荐
相关产品推荐

