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

如何在Python FastAPI中返回流数据响应并在前端使用?

FastAPI流式响应返回与前端使用解决方案

核心问题诊断

你的代码存在两个关键问题导致无法返回流式响应:

  1. read_stream是异步函数,直接await read_stream()会等待其执行完毕,但内部的while True是无限循环,永远不会触发返回。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 11:11:18