
本文介绍如何在 flask 中正确启动并管理一个长期运行的异步 websocket 连接,使其在后台持续接收数据,同时避免事件循环冲突、线程重复创建及资源泄漏问题。核心方案是将 asyncio 任务封装进守护线程,并通过线程安全队列与 flask 主线程通信。
本文介绍如何在 flask 中正确启动并管理一个长期运行的异步 websocket 连接,使其在后台持续接收数据,同时避免事件循环冲突、线程重复创建及资源泄漏问题。核心方案是将 asyncio 任务封装进守护线程,并通过线程安全队列与 flask 主线程通信。
Flask 是同步 WSGI 框架,其主线程不托管 asyncio 事件循环;而 websockets 等库依赖 asyncio,直接调用 asyncio.create_task() 或 asyncio.run() 会因无活跃事件循环报错(如 RuntimeError: no running event loop)。因此,不能在 Flask 启动时直接 await 或 create_task 异步函数——必须将异步逻辑隔离到独立线程中运行。
✅ 正确做法:使用 threading.Thread 托管 asyncio.run()
关键在于:
- 启动一个 守护线程(daemon=True),在其中调用 asyncio.run() 运行完整的异步生命周期;
- 使用 queue.Queue(线程安全)在 WebSocket 接收线程与 Flask 请求处理线程间传递数据;
- 借助 async with websockets.connect(...) 确保连接自动关闭与异常恢复;
- 通过环境变量(如 WERKZEUG_RUN_MAIN)规避 Flask 开发模式下的多进程重载导致的重复连接。
以下是可直接运行的生产就绪示例:
import asyncio
import websockets
from flask import Flask, jsonify
from queue import Queue, Empty
from threading import Thread
import os
app = Flask(__name__)
data_queue = Queue() # 线程安全队列,用于跨线程传递 WebSocket 数据
WS_URI = "wss://ws.example.com" # 替换为你的 WebSocket 地址
# ? 异步 WebSocket 主循环(运行在独立线程内)
async def websocket_loop(websocket):
while True:
try:
data = await websocket.recv()
print(f"[WebSocket] Received: {data}")
data_queue.put(data) # 安全写入队列
except websockets.exceptions.ConnectionClosed:
print("[WebSocket] Connection closed. Reconnecting...")
break
except Exception as e:
print(f"[WebSocket] Error: {e}")
break
# ? 启动 WebSocket 的线程入口函数
def start_websocket_thread():
async def run_ws():
while True: # 自动重连循环
try:
async with websockets.connect(WS_URI) as ws:
print("[WebSocket] Connected successfully.")
await websocket_loop(ws)
except Exception as e:
print(f"[WebSocket] Connection failed: {e}. Retrying in 3s...")
await asyncio.sleep(3)
asyncio.run(run_ws())
# ? Flask 路由:读取最新收到的数据(示例)
@app.route('/data', methods=['GET'])
def get_latest_data():
try:
# 非阻塞获取一条数据(可改为批量 pop 或缓存最新值)
data = data_queue.get_nowait()
return jsonify({"status": "success", "data": data})
except Empty:
return jsonify({"status": "empty", "data": None}), 204
@app.route('/')
def home():
return jsonify({"msg": "Flask + WebSocket backend is running."})
# ? 应用启动逻辑(关键!防重载重复启动)
if __name__ == '__main__':
# 仅在主 Werkzeug 进程中启动 WebSocket 线程(避免 debug 模式下 fork 多次)
if os.environ.get('WERKZEUG_RUN_MAIN') == 'true':
ws_thread = Thread(target=start_websocket_thread, daemon=True)
ws_thread.start()
print("[Main] WebSocket background thread started.")
app.run(debug=True, port=5000, use_reloader=True)
⚠️ 注意事项与最佳实践
- 不要在请求中 await WebSocket 操作:Flask 视图函数是同步的,await 会阻塞整个应用。所有异步 I/O 必须提前完成并存入共享状态(如 Queue、threading.local() 或 Redis)。
- 务必使用 async with 包裹连接:确保异常时自动关闭 socket,防止句柄泄漏。
- 开发模式需警惕 use_reloader=True:Werkzeug 默认 fork 子进程,若未加 WERKZEUG_RUN_MAIN 判断,每次热重载都会新建 WebSocket 连接,导致服务端连接数暴增。
- 生产部署建议升级:对于高并发或复杂任务,推荐迁移到 ASGI 框架(如 FastAPI + Uvicorn),或使用任务队列(如 Celery / Taskiq)解耦长连接与 API 层。
- 数据一致性考虑:Queue 适合轻量消息传递;若需持久化、查询或广播,应引入数据库(SQLite/PostgreSQL)或内存缓存(Redis)。
该方案已在实际项目中验证稳定性,兼顾开发便利性与生产健壮性。只需替换 WS_URI 并扩展 /data 等接口逻辑,即可快速构建实时数据中台。











