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

如何将OpenAI流式回复存储到RAG聊天机器人的聊天历史数据库?

问题解决:流式AI回复存入聊天历史的异步处理修复

核心问题

你遇到的问题本质是异步流式生成器的执行时机偏差:finalize_chat_history在StreamingResponse返回客户端前就被调用,但此时event_generator还未开始迭代流式chunk,导致complete_response为空,无法保存完整AI回复。

解决方案

将聊天历史的保存逻辑移到流式生成器迭代完成后执行,也就是把finalize_chat_history的调用放到event_generator内部,在所有chunk处理完毕(或异常结束)后触发。

修改后的代码

API路由代码(关键修改部分)

@router.post("/message_stream")
async def message_stream(request: Request, message: Message):
    try:
        history_key = "public_chat_history"
        # 获取现有聊天历史
        chat_history = chat_service.get_chat_history(request, history_key=history_key)
        # 保存用户消息到历史
        chat_service.add_message_to_history(
            request,
            history_key=history_key,
            message=message.text,
            role="user"
        )
    
        complete_response = []  # 用于收集流式返回的所有chunk
    
        # 流式响应生成器
        async def event_generator():
            try:
                async for chunk in rag_service.arag_stream(
                    input=message.text,
                    chat_history=chat_history,
                    authorization="public"
                ):
                    if chunk:
                        complete_response.append(chunk)
                        print(complete_response)
                        yield chunk  # 实时向客户端返回chunk
            except Exception as e:
                logging.error(f"流式响应出错: {e}")
                raise e
            finally:
                # 流式迭代完成/异常结束后,执行聊天历史保存
                await finalize_chat_history()
    
        # 保存完整AI回复到聊天历史的方法
        async def finalize_chat_history():
            try:
                print("进入finalize_chat_history")
                print(complete_response)
                full_response = ''.join(complete_response)
                if full_response:
                    chat_service.add_message_to_history(
                        request,
                        history_key=history_key,
                        message=full_response,
                        role="ai"
                    )
                    logging.info(f"完整回复已保存到聊天历史: {full_response}")
                else:
                    logging.warning("完整回复为空,跳过保存")
            except Exception as e:
                logging.error(f"保存聊天历史出错: {e}")
    
        # 创建流式响应对象
        response = StreamingResponse(
            event_generator(),
            media_type="text/event-stream",
            headers={
                "Cache-Control": "no-cache",
                "Connection": "keep-alive",
                "Content-Type": "text/event-stream",
                "X-Accel-Buffering": "no"  # 禁用代理缓冲
            }
        )
    
        return response
    
    except Exception as e:
        logging.error(f"流式接口整体出错: {e}")
        raise HTTPException(status_code=500, detail=str(e))

修改说明

  1. 调整保存逻辑的执行时机:将await finalize_chat_history()移到event_generator的finally块中,确保只有在所有流式chunk迭代完成(或异常终止)后才执行保存操作,此时complete_response已收集完整的AI回复。
  2. 保留原有收集逻辑:维持complete_response逐chunk拼接的逻辑,确保能生成完整的AI回复内容。

内容的提问来源于stack exchange,提问作者Abdullah Muhammad Moosa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 13:52:08