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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 13:48:21