如何保障OpenAI Chat Completion API流式传输不中断及交付后断开连接
FastAPI OpenAI Chat Completion接口问题解决方案
问题1:流式会话被非流式请求中断
核心原因
流式接口若使用同步阻塞逻辑或共享全局非线程安全资源,会占用FastAPI的事件循环,导致新的非流式请求抢占资源,中断正在进行的流式传输。
解决方案
强制异步处理流式任务
所有流式生成逻辑必须异步化,避免阻塞事件循环。若调用的LLM SDK为同步版本,用线程池封装为异步操作:from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio from concurrent.futures import ThreadPoolExecutor import openai app = FastAPI() # 初始化线程池,隔离同步任务 executor = ThreadPoolExecutor(max_workers=5) async def sync_to_async(func, *args, **kwargs): loop = asyncio.get_event_loop() return await loop.run_in_executor(executor, func, *args, **kwargs) async def stream_generator(messages): # 将同步流式调用封装为异步 completion = await sync_to_async( openai.ChatCompletion.create, model="gpt-3.5-turbo", messages=messages, stream=True ) for chunk in completion: yield f"data: {chunk.json()}\n\n" @app.post("/chat/completion/stream") async def chat_stream(messages: list[dict]): # 每个请求生成独立的异步生成器,隔离资源 return StreamingResponse(stream_generator(messages), media_type="text/event-stream")隔离请求级资源
禁止在不同请求间共享LLM客户端实例、会话状态等非线程安全对象,通过依赖注入为每个请求创建独立实例。
问题2:流式API无请求完成日志且连接未断开
核心原因
流式响应默认在生成器耗尽后关闭连接,但缺少明确的结束标记和日志回调,导致客户端可能维持连接,且FastAPI未打印请求完成日志。
解决方案
添加流式结束标记
在生成器末尾输出标准结束信号,告知客户端会话完成:async def stream_generator(messages): completion = await sync_to_async( openai.ChatCompletion.create, model="gpt-3.5-turbo", messages=messages, stream=True ) for chunk in completion: yield f"data: {chunk.json()}\n\n" # 输出结束标记,触发客户端断开连接 yield "data: [DONE]\n\n"自定义中间件打印完成日志
通过中间件包装流式响应的生成器,在生成器耗尽后打印与非流式一致的请求日志:from starlette.middleware.base import BaseHTTPMiddleware from starlette.requests import Request from starlette.responses import StreamingResponse import logging logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) class StreamingLogMiddleware(BaseHTTPMiddleware): async def dispatch(self, request: Request, call_next): response = await call_next(request) if isinstance(response, StreamingResponse): original_generator = response.body_iterator async def wrapped_generator(): try: async for chunk in original_generator: yield chunk finally: # 打印请求完成日志 logger.info(f'{request.client.host}:{request.client.port} - "{request.method} {request.url.path}?{request.query_params}" 200 OK') response.body_iterator = wrapped_generator() return response app.add_middleware(StreamingLogMiddleware)
内容的提问来源于stack exchange,提问作者James K J
相关产品推荐
相关产品推荐

