在Flask应用后台运行异步WebSocket循环报错求助
Flask + WebSocket 协程错误解决方案
错误原因
你的代码出现两个问题的核心原因:
- Flask 的
app.run()是同步阻塞调用,会直接占据主线程,导致你创建的 asyncio 任务完全没机会运行 asyncio.create_task()必须在已启动的事件循环中调用,你调用它时事件循环还未启动,因此抛出RuntimeError: no running event loop- 协程
websocket_loop从未被 await,所以触发RuntimeWarning: coroutine 'websocket_loop' was never awaited
解决方案一:Flask 结合线程运行事件循环
通过单独线程启动 asyncio 事件循环,避免被 Flask 的同步运行逻辑阻塞,同时处理同步视图调用异步方法的问题:
import asyncio import websockets from flask import Flask, jsonify, request from threading import Thread app = Flask(__name__) websocket = None # 存储WebSocket收到的流数据,供API调用 received_data = [] async def websocket_setup(): global websocket uri = 'wss://ws.example.com' websocket = await websockets.connect(uri) print("WebSocket connected") async def websocket_loop(): global websocket, received_data await websocket_setup() try: while True: data = await websocket.recv() print("Received:", data) received_data.append(data) # 限制存储数量,防止内存溢出 if len(received_data) > 100: received_data.pop(0) except websockets.exceptions.ConnectionClosed: print("WebSocket连接断开,正在重连...") await websocket_setup() def start_async_loop(): # 创建并运行独立的事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) loop.run_until_complete(websocket_loop()) @app.get('/') def home(): return jsonify({'msg':'hello'}) @app.get('/latest-data') def get_latest_data(): # 返回最新的一条WebSocket数据 return jsonify({'latest_data': received_data[-1] if received_data else None}) @app.post('/send-message') def send_message(): # 在同步视图中调用异步WebSocket发送方法 if not websocket: return jsonify({'error': 'WebSocket未连接'}), 500 message = request.json.get('message') if not message: return jsonify({'error': '未提供消息内容'}), 400 loop = asyncio.get_event_loop() asyncio.run_coroutine_threadsafe(websocket.send(message), loop) return jsonify({'status': '消息已发送'}) if __name__ == '__main__': # 启动后台线程运行事件循环 async_thread = Thread(target=start_async_loop, daemon=True) async_thread.start() app.run(debug=False)
解决方案二:改用 FastAPI(原生支持异步)
FastAPI 天生支持异步逻辑,无需额外线程处理,是更适配的方案:
import asyncio import websockets from fastapi import FastAPI from pydantic import BaseModel app = FastAPI() websocket = None received_data = [] async def websocket_loop(): global websocket, received_data uri = 'wss://ws.example.com' while True: try: websocket = await websockets.connect(uri) print("WebSocket connected") while True: data = await websocket.recv() print("Received:", data) received_data.append(data) if len(received_data) > 100: received_data.pop(0) except websockets.exceptions.ConnectionClosed: print("WebSocket连接断开,5秒后重连...") await asyncio.sleep(5) except Exception as e: print(f"WebSocket错误: {str(e)}") await asyncio.sleep(5) # FastAPI启动时自动启动WebSocket后台任务 @app.on_event("startup") async def startup_event(): asyncio.create_task(websocket_loop()) @app.get("/") async def home(): return {"msg": "hello"} @app.get("/latest-data") async def get_latest_data(): return {"latest_data": received_data[-1] if received_data else None} class Message(BaseModel): message: str @app.post("/send-message") async def send_message(msg: Message): if not websocket: return {"error": "WebSocket未连接"}, 500 await websocket.send(msg.message) return {"status": "消息已发送"}
内容的提问来源于stack exchange,提问作者Rishabh Gupta
相关产品推荐
相关产品推荐

