FastAPI中返回StreamingResponse后如何执行数据库插入操作?
FastAPI流式返回同时保存完整文本到数据库
核心思路
在生成流式响应片段的过程中收集完整文本,待所有片段生成完毕后,通过后台任务异步将完整文本插入数据库,既不阻塞流式返回,也能确保数据持久化。
具体实现代码
1. 导入所需依赖
from fastapi import FastAPI, BackgroundTasks, Depends from sqlalchemy.orm import Session from your_module import Input, generate_streaming_response, get_db, YourModel # 替换为实际业务模块
2. 定义数据库保存函数
def save_full_text_to_db(text: str): # 独立创建数据库会话,避免依赖请求会话的生命周期 db = next(get_db()) try: # 创建数据库模型实例并执行插入 db_entry = YourModel(content=text) db.add(db_entry) db.commit() db.refresh(db_entry) finally: db.close()
3. 修改路由函数
@app.post("/foo_stream") def foo_stream(input: Input, background_tasks: BackgroundTasks): full_text = [] def streaming_generator(): # 遍历流式生成的片段,同时收集完整文本 for chunk in generate_streaming_response(input): full_text.append(chunk) yield chunk # 所有片段生成完成后,添加后台任务保存完整文本 background_tasks.add_task(save_full_text_to_db, ''.join(full_text)) return StreamingResponse(streaming_generator(), media_type="text/plain")
关键点说明
- 流式生成器封装:通过自定义生成器,在返回每个文本片段的同时,将片段追加到
full_text列表中,确保完整文本的收集。 - 后台任务异步执行:使用FastAPI的
BackgroundTasks,在流式响应完成后异步触发数据库插入操作,不会阻塞客户端接收流式内容。 - 独立数据库会话:在保存函数中自行创建数据库会话,避免请求结束后会话被关闭导致的数据库操作失败。
异步场景适配(如果生成器是异步的)
如果generate_streaming_response是异步生成器,只需调整为异步写法:
@app.post("/foo_stream") async def foo_stream(input: Input, background_tasks: BackgroundTasks): full_text = [] async def streaming_generator(): async for chunk in generate_streaming_response(input): full_text.append(chunk) yield chunk background_tasks.add_task(save_full_text_to_db, ''.join(full_text)) return StreamingResponse(streaming_generator(), media_type="text/plain")
内容的提问来源于stack exchange,提问作者gameveloster
相关产品推荐
相关产品推荐

