Quart多WebSocket使用时第二个连接延迟运行问题及聊天服务示例咨询
Quart多WebSocket阻塞问题解决及聊天服务器示例
问题成因
- 核心原因是你使用了同步的SQLite操作接口,同步IO操作会直接阻塞整个asyncio事件循环,在数据库查询返回前,事件循环无法分配资源处理
/printar路由的WebSocket连接请求,就会出现第二个连接必须等第一个的数据库操作完成才能启动的现象。 - 你现有
/mensagens路由的异步函数存在语法结构错误:await asyncio.sleep(1)放在了try块外部,会直接触发异常终止函数运行,且没有循环结构无法持续推送新消息。
修复步骤
- 替换同步SQLite驱动为异步的
aiosqlite,所有数据库操作都通过await调用,不会阻塞事件循环 - 调整WebSocket路由的逻辑结构,添加循环保证长连接持续运行
修复后核心代码示例
# 安装依赖:pip install quart aiosqlite from quart import Quart, websocket import asyncio import aiosqlite app = Quart(__name__) # 会话消息推送WebSocket @app.websocket('/mensagens/<dialog_id>') async def mensagens(dialog_id): # 连接异步sqlite async with aiosqlite.connect('你的数据库路径.db') as db: while True: try: # 异步查询数据库新消息,不会阻塞事件循环 cursor = await db.execute("SELECT content FROM messages WHERE dialog_id = ?", (dialog_id,)) output = await cursor.fetchone() if output: await websocket.send(f"{output[0]}") # 每秒查询一次 await asyncio.sleep(1) except Exception as e: print("消息推送异常:", e) break # 消息接收WebSocket @app.websocket('/printar/<dialog_id>') async def printar(dialog_id): async with aiosqlite.connect('你的数据库路径.db') as db: try: while True: data = await websocket.receive() # 异步写入数据库 await db.execute("INSERT INTO messages (dialog_id, content) VALUES (?, ?)", (dialog_id, data)) await db.commit() print(f"收到对话{dialog_id}的消息:{data}") except Exception as e: print("消息接收异常:", e) if __name__ == "__main__": try: app.run() except KeyboardInterrupt: print("=====\nAdeus!\n=====") except Exception as e: print(e)
简单Quart聊天室完整实现
from quart import Quart, render_template_string, websocket from collections import defaultdict import asyncio app = Quart(__name__) # 按房间存储连接对象 room_connections = defaultdict(set) # 前端页面 @app.route('/<room_id>') async def index(room_id): return await render_template_string(''' <!DOCTYPE html> <html> <body> <ul id="messages"></ul> <input type="text" id="msgInput" placeholder="输入消息"> <button onclick="sendMsg()">发送</button> <script> // 接收消息的ws const recvWs = new WebSocket(`ws://${window.location.host}/recv/{{room_id}}`); recvWs.onmessage = (event) => { const li = document.createElement('li'); li.textContent = event.data; document.getElementById('messages').appendChild(li); } // 发送消息的ws const sendWs = new WebSocket(`ws://${window.location.host}/send/{{room_id}}`); function sendMsg() { const input = document.getElementById('msgInput'); sendWs.send(input.value); input.value = ''; } </script> </body> </html> ''', room_id=room_id) # 消息接收路由 @app.websocket('/send/<room_id>') async def send(room_id): while True: data = await websocket.receive() # 推送给同房间所有连接 for conn in room_connections[room_id]: await conn.send(data) # 消息订阅路由 @app.websocket('/recv/<room_id>') async def recv(room_id): room_connections[room_id].add(websocket._get_current_object()) try: while True: await asyncio.sleep(3600) # 保持连接存活 finally: room_connections[room_id].remove(websocket._get_current_object()) if __name__ == "__main__": app.run()
内容的提问来源于stack exchange,提问作者Joao Pedro Lourenco Affonso
相关产品推荐
相关产品推荐

