FastAPI中如何将WebSocket数据传入/chart-data的StreamingResponse?
解决方案:将WebSocket接收的温度数据接入SSE图表接口
要实现WebSocket接收树莓派数据后,通过/chart-data的SSE接口推送给前端图表,核心是用异步队列在两个独立的请求处理逻辑间传递数据。以下是具体实现步骤及修正后的代码:
1. 全局异步队列(数据中转站)
首先定义一个线程/协程安全的异步队列,用于存储WebSocket收到的温度数据,让WebSocket端点和SSE端点能共享数据:
import asyncio from fastapi import FastAPI, WebSocket, Request from fastapi.responses import StreamingResponse import json from datetime import datetime import logging logger = logging.getLogger(__name__) application = FastAPI() # 全局异步队列,限制队列大小避免内存溢出 data_queue = asyncio.Queue(maxsize=10)
2. 修正WebSocket端点(接收并转发数据到队列)
修改原WebSocket处理函数,解析树莓派传来的温度数据,格式化后存入队列:
@application.websocket("/ws") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() # 获取客户端真实IP(适配Heroku反向代理) client_ip = websocket.headers.get("X-Forwarded-For", websocket.client.host) logger.info("Client %s connected", client_ip) try: while True: received = await websocket.receive() logger.info(received) # 解析JSON数据,容错处理缺失字段 msg = json.loads(received["text"]) temperature = msg.get("temperature") if temperature is None: await websocket.send_text("Error: missing 'temperature' field") continue # 构造SSE所需的JSON格式数据 sse_data = json.dumps({ "time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "value": temperature }) # 队列满时自动丢弃旧数据,避免阻塞 if data_queue.full(): await data_queue.get() await data_queue.put(sse_data) await websocket.send_text("data received") except Exception as e: logger.error("Connection terminated: %s", str(e)) print(f"Terminated: {e}") finally: await websocket.close()
3. 重写SSE接口(从队列取数据并推送给前端)
编写专门的SSE数据生成器,从队列中持续获取数据,按照SSE规范格式推送给前端:
async def sse_data_generator(request: Request): while True: # 检查客户端是否断开连接,避免无效推送 if await request.is_disconnected(): break # 从队列取数据,超时后发送心跳维持连接 try: data = await asyncio.wait_for(data_queue.get(), timeout=5) except asyncio.TimeoutError: # 心跳数据(前端可忽略空值) data = json.dumps({"time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "value": None}) # 严格遵循SSE格式:data: {JSON}\n\n yield f"data: {data}\n\n" @application.get("/chart-data") async def chart_data(request: Request) -> StreamingResponse: response = StreamingResponse(sse_data_generator(request), media_type="text/event-stream") response.headers["Cache-Control"] = "no-cache" response.headers["X-Accel-Buffering"] = "no" response.headers["Connection"] = "keep-alive" return response
关键说明
- 异步队列的作用:FastAPI中WebSocket和SSE是独立的请求协程,通过
asyncio.Queue可以安全地在协程间传递数据,避免线程安全问题。 - SSE格式要求:前端
EventSource只能识别data: {内容}\n\n格式的消息,必须严格遵循,否则无法解析。 - 容错处理:添加了字段缺失检查、队列满时的旧数据丢弃、客户端断开检测,避免服务端异常阻塞。
- 原代码错误修正:
- 原
generate_random_data函数混淆了WebSocket端点和数据生成器的职责,拆分后职责更清晰。 - 修正了
yieldy函数引用外部变量的错误,改为直接使用传入参数。
- 原
内容的提问来源于stack exchange,提问作者aj91
相关产品推荐
相关产品推荐

