FastAPI EventSourceResponse SSE数据丢失问题排查与解决求助
问题描述
我使用FastAPI结合sse_starlette.sse的EventSourceResponse实现ChatGPT输出流式传输,初期基于sse.js的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(); };
问题原因及解决办法
核心原因
- 空chunk未过滤:OpenAI流式响应中部分
chunk的delta字段不含content(比如仅返回角色信息或结束标记),后端仍生成包含空content的JSON并发送,导致前端无法识别有效数据。 - SSE格式不规范:
sse_starlette的EventSourceResponse要求每个事件遵循SSE标准格式(以data:开头,结尾加两个换行符),直接yield JSON字符串不符合规范,引发前端解析异常。
解决步骤
后端代码调整
- 过滤空内容并规范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"
前端代码优化
- 增加异常处理,避免解析失败中断逻辑:
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
相关产品推荐
相关产品推荐

