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

FastAPI EventSourceResponse SSE数据丢失问题排查与解决求助

问题描述

我使用FastAPI结合sse_starlette.sse的EventSourceResponse实现ChatGPT输出流式传输,初期基于sse.js的UI能正常接收响应,但后续仅能收到事件却无法获取实际数据。以下是相关代码与UI日志截图:

UI端日志截图

UI日志截图

FastAPI后端代码(Python)

from sse_starlette.sse import EventSourceResponse

@router.post("/stream")
async def chat_stream(chat_msg: ChatMessage, user: CurrentUser = Depends(get_current_user)):
    try:
        return EventSourceResponse(chat_stram_reply_with_history(chat_msg, curr_user=user))
    except ValueError as error:
        log.error(error)
        raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail=error.args)

async def chat_stram_reply_with_history(chat_msg: chat.ChatMessage, curr_user: user.CurrentUser):
    try:
        response = openai.ChatCompletion.create(
            model=chat_msg.model,
            messages=msgs,
            temperature=chat_msg.temperature,
            max_tokens=chat_msg.max_tokens,
            top_p=chat_msg.top_p,
            stream=True,
            user=str(curr_user.sub)
        )
        for chunk in response:
            collected_chunks.append(chunk)
            choice = chunk['choices'][0]
            chunk_message = choice['delta']
            collected_messages.append(chunk_message)
            res = {
                "message_id": reply_msg_id_str,
                "content": chunk_message.get('content', ''),
                "conversation_id": conversation_id_str,
                "parent_message_id": reply_parent_message_id_str,
                "finish_reason": choice['finish_reason']
            }
            log.debug(f"User: {curr_user.preferred_username} | Response: {res}") # 日志能正常打印,但UI端数据为空
            yield json.dumps(res)
    except openai.error.APIError as e:
        error_msg = f"OpenAI API returned an API Error: {e}"
    except openai.error.APIConnectionError as e:
        error_msg = f"Failed to connect to OpenAI API: {e}"
    except openai.error.RateLimitError as e:
        error_msg = f"OpenAI API request exceeded rate limit: {e}"
    except Exception as e:
        error_msg = f"OpenAI API returned an Exception: {e}"
    final_msg = {
       'finish_reason': 'done'
    }
    yield json.dumps(final_msg)

Vue3前端代码

import { SSE } from 'sse'
const sendMessage = () => {
  const source = new SSE(`${aiUrl}/api/v1/chat/stream`, {
    headers: {
      "Content-Type": "application/json",
      Authorization: getToken(),
    },
    method: "POST",
    withCredentials: true,
    payload: JSON.stringify(payload),
  });
  source.addEventListener("message", (e) => {
    console.log("res", e);
    console.log("res", e.data);
    const resp = JSON.parse(e.data || {});
    if (resp.finish_reason === "done") {
      source.close();
    }
  });
  source.stream();
};

问题原因及解决办法

核心原因

  1. 空chunk未过滤:OpenAI流式响应中部分chunk的delta字段不含content(比如仅返回角色信息或结束标记),后端仍生成包含空content的JSON并发送,导致前端无法识别有效数据。
  2. SSE格式不规范:sse_starlette的EventSourceResponse要求每个事件遵循SSE标准格式(以data: 开头,结尾加两个换行符),直接yield JSON字符串不符合规范,引发前端解析异常。

解决步骤

后端代码调整

  1. 过滤空内容并规范SSE格式:仅当有有效内容或结束标记时,按标准格式发送数据:
async def chat_stram_reply_with_history(chat_msg: chat.ChatMessage, curr_user: user.CurrentUser):
    # 改为局部变量,避免多用户数据混淆
    collected_chunks = []
    collected_messages = []
    try:
        response = openai.ChatCompletion.create(
            model=chat_msg.model,
            messages=msgs,
            temperature=chat_msg.temperature,
            max_tokens=chat_msg.max_tokens,
            top_p=chat_msg.top_p,
            stream=True,
            user=str(curr_user.sub)
        )
        for chunk in response:
            choice = chunk['choices'][0]
            chunk_message = choice['delta']
            collected_chunks.append(chunk)
            collected_messages.append(chunk_message)
            
            # 仅发送有内容的消息
            content = chunk_message.get('content', '')
            if content:
                res = {
                    "message_id": reply_msg_id_str,
                    "content": content,
                    "conversation_id": conversation_id_str,
                    "parent_message_id": reply_parent_message_id_str,
                    "finish_reason": None
                }
                # 按SSE标准格式包装
                yield f"data: {json.dumps(res)}\n\n"
            # 发送结束标记
            elif choice['finish_reason']:
                res = {
                    "finish_reason": choice['finish_reason']
                }
                yield f"data: {json.dumps(res)}\n\n"
    except Exception as e:
        error_msg = f"请求异常: {str(e)}"
        # 发送错误信息给前端
        yield f"data: {json.dumps({'error': error_msg, 'finish_reason': 'error'})}\n\n"
    # 最终结束信号
    yield f"data: {json.dumps({'finish_reason': 'done'})}\n\n"

前端代码优化

  1. 增加异常处理,避免解析失败中断逻辑:
import { SSE } from 'sse'
const sendMessage = () => {
  const source = new SSE(`${aiUrl}/api/v1/chat/stream`, {
    headers: {
      "Content-Type": "application/json",
      Authorization: getToken(),
    },
    method: "POST",
    withCredentials: true,
    payload: JSON.stringify(payload),
  });
  source.addEventListener("message", (e) => {
    console.log("res", e);
    console.log("res", e.data);
    let resp = {};
    try {
      if (e.data) {
        resp = JSON.parse(e.data);
      }
    } catch (err) {
      console.error("解析SSE数据失败:", err);
      return;
    }
    if (resp.finish_reason === "done" || resp.finish_reason) {
      source.close();
    } else if (resp.content) {
      // 处理有效内容,比如追加到聊天框
      console.log("收到内容:", resp.content);
    }
  });
  // 监听错误事件
  source.addEventListener("error", (err) => {
    console.error("SSE连接错误:", err);
    source.close();
  });
  source.stream();
};

内容的提问来源于stack exchange,提问作者Sasank M

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 14:24:51