WebSocket无提示自动关闭问题:推送客户端余额更新时发送accepted消息后无法传输数据
问题分析与解决方案
看起来你的WebSocket连接在发送accepted消息后就自动关闭,核心原因是FastAPI的WebSocket端点函数执行完毕后,会自动终止连接。当前你的api_websocket_wait_update函数在发送accepted后没有任何阻塞逻辑,函数直接结束,导致连接被关闭。此外还有几处代码细节问题需要修正,我来一步步帮你解决:
核心问题与修正步骤
1. 让WebSocket端点保持运行
FastAPI的WebSocket端点需要函数持续运行才能维持连接。我们可以用asyncio.Event来等待余额更新事件,直到回调触发或者客户端主动断开。
2. 修正Token不存在的错误处理
当没有获取到access_token时,返回JSONResponse是无效的——因为WebSocket已经建立了双向连接,应该发送错误消息后主动关闭连接。
3. 修复Firebase监听器的回调参数错误
你的get_contract_listener中,on_snapshot的回调参数处理有误:单个文档的on_snapshot回调第一个参数就是DocumentSnapshot对象,不需要取doc[0],这会导致无法正确获取余额数据。
4. 正确处理异步任务提交
在同步的on_snapshot回调中执行异步函数,不能用asyncio.run(会创建新事件循环),应该用asyncio.create_task将异步任务提交到当前FastAPI的事件循环中。
修改后的完整代码
后端Python代码
import asyncio from fastapi import WebSocket, WebSocketDisconnect from google.cloud.firestore import DocumentSnapshot, Watch, FieldFilter from typing import Callable, Awaitable, cast # 修正get_contract_listener函数 def get_contract_listener(access_token: str, callback: Callable[[int | None], Awaitable[None]]) -> None | Watch: logger.debug('get contract listener token=%s', access_token) contract = db.collection('Contract').where(filter=FieldFilter('access_token', '==', access_token)).get() if not contract: logger.debug('contract not found') return async def wrap(document: DocumentSnapshot): await callback((document.to_dict() or {}).get('last_saldo')) # 修正回调参数,用create_task提交异步任务 return cast(DocumentReference, contract[0].reference).on_snapshot( lambda doc, _, __: asyncio.create_task(wrap(doc)) ) # 修正WebSocket端点函数 @router.websocket('/wait_update') async def api_websocket_wait_update(socket: WebSocket): await socket.accept() token = socket.cookies.get('access_token') # 修正Token不存在的处理逻辑 if token is None: await socket.send_json({'status': 'fail', 'detail': 'no access token'}) await socket.close(1008) # 1008表示权限相关的关闭码 return update_event = asyncio.Event() listener = None async def callback(balance: int | None) -> None: nonlocal listener # 先取消监听器,避免重复触发 if listener is not None: listener.unsubscribe() # 发送更新消息 await socket.send_json({'message': 'updated', 'balance': balance}) # 触发事件,让主函数退出,关闭连接 update_event.set() listener = get_contract_listener(token, callback) if listener is None: await socket.send_json({'message': 'failed', 'detail': 'contract not found'}) await socket.close(1007) return await socket.send_json({'message': 'accepted'}) try: # 等待更新事件或者客户端断开连接 await asyncio.wait_for(update_event.wait(), timeout=None) except WebSocketDisconnect: # 客户端主动断开时,清理监听器 if listener is not None: listener.unsubscribe() logger.debug('Client disconnected') finally: # 确保连接关闭 await socket.close()
前端JS代码(添加错误提示优化)
function create_socket() { const socket = new WebSocket('/api/contract/wait_update'); socket.addEventListener('open', () => { console.log('open websocket connection'); }); socket.addEventListener('message', (e) => { let data; try { data = JSON.parse(e.data); } catch (err) { console.error('Failed to parse websocket message:', err); return; } if (data.message == 'failed') { console.error(`error from websocket: ${data.detail}`); socket.close(); return; } console.log(`message from websocket: ${e.data}`); if (data.message == 'updated') { console.log('balance updated'); // 执行你的余额更新逻辑... socket.close(); // 更新完成后主动关闭连接 } }); socket.addEventListener('close', (event) => { console.log(`close websocket connection, code: ${event.code}, reason: ${event.reason}`); }); socket.addEventListener('error', (err) => { console.error('WebSocket error:', err); }); }
关键修改说明
- 用
asyncio.Event维持连接:端点函数会一直等待update_event被触发(余额更新时),避免函数提前结束导致连接关闭。 - 修正Firebase监听器参数:将
doc[0]改为doc,确保正确传递DocumentSnapshot对象到回调函数。 - 异步任务提交:用
asyncio.create_task替代asyncio.run,避免在现有事件循环中创建新循环引发冲突。 - 完善错误处理:Token不存在时发送WebSocket错误消息并关闭连接,客户端断开时清理Firebase监听器,避免资源泄漏。
内容的提问来源于stack exchange,提问作者firedotguy
相关产品推荐
相关产品推荐

