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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 03:15:03