如何在Python FastAPI中返回流数据响应并在前端使用?
FastAPI流式响应返回与前端使用解决方案
核心问题诊断
你的代码存在两个关键问题导致无法返回流式响应:
read_stream是异步函数,直接await read_stream()会等待其执行完毕,但内部的while True是无限循环,永远不会触发返回。print_message仅打印数据,没有将消息输出到响应流中,无法传递给前端。
修正后的后端代码
通过异步生成器配合消息队列中转,实现流式数据返回:
import asyncio import json from fastapi import FastAPI, StreamingResponse import schwab app = FastAPI() token_path = "your_token_path" API_KEY = "your_api_key" API_SECRET = "your_api_secret" ACCOUNT_ID = "your_account_id" @app.get("/streaming/{symbol}") async def getStreamingQuoteBySymbol(symbol): try: client = schwab.auth.client_from_token_file(token_path=token_path, api_key=API_KEY, app_secret=API_SECRET) streamClient = schwab.streaming.StreamClient(client=client, account_id=ACCOUNT_ID) # 用队列中转流式消息,解决同步处理器与异步生成器的衔接问题 message_queue = asyncio.Queue() def message_handler(message): # 序列化消息并加入队列,换行符方便前端逐行解析 json_data = json.dumps(message) + "\n" asyncio.create_task(message_queue.put(json_data)) async def stream_generator(): try: await streamClient.login() print("Login successfully") except Exception as e: print("Failed to login streaming api.") yield json.dumps({"error": "登录流式API失败"}) + "\n" return # 绑定消息处理器 streamClient.add_level_one_equity_handler(message_handler) await streamClient.level_one_equity_subs([symbol]) # 持续从队列取消息并输出到响应流 while True: data = await message_queue.get() yield data # 异步处理streamClient消息,维持连接活跃 asyncio.create_task(streamClient.handle_message()) # 指定媒体类型为SSE,适配前端EventSource return StreamingResponse(stream_generator(), media_type="text/event-stream") except Exception as e: print("Error message: ", e) return {"status": "流式请求失败"}
关键修改说明
- 使用
asyncio.Queue中转消息:由于schwab的流处理器是同步函数,无法直接在内部yield,通过队列实现同步到异步的消息传递。 - 异步生成器
stream_generator:通过yield逐次输出消息,FastAPI的StreamingResponse会自动处理流式返回。 - 媒体类型选择
text/event-stream:这是服务器发送事件(SSE)的标准类型,前端可以用EventSource轻松接收。
前端接收示例
方案1:使用EventSource(推荐,SSE原生支持)
const targetSymbol = "AAPL"; // 替换为目标股票代码 const eventSource = new EventSource(`/streaming/${targetSymbol}`); // 接收流式消息 eventSource.onmessage = (event) => { const streamData = JSON.parse(event.data); console.log("实时报价数据:", streamData); // 这里添加UI更新逻辑,比如渲染到页面 }; // 处理错误 eventSource.onerror = (error) => { console.error("流连接出错:", error); eventSource.close(); };
方案2:使用Fetch API + ReadableStream
async function fetchRealTimeQuote(symbol) { const response = await fetch(`/streaming/${symbol}`); const reader = response.body.getReader(); const decoder = new TextDecoder(); while (true) { const { done, value } = await reader.read(); if (done) break; const rawText = decoder.decode(value); // 按换行分割多条消息,过滤空内容 const messages = rawText.split("\n").filter(msg => msg.trim()); messages.forEach(msg => { const data = JSON.parse(msg); console.log("收到数据:", data); // 执行UI更新操作 }); } } // 调用示例 fetchRealTimeQuote("AAPL");
内容的提问来源于stack exchange,提问作者Adnan Rahman
相关产品推荐
相关产品推荐

