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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 10:17:58