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

如何保障OpenAI Chat Completion API流式传输不中断及交付后断开连接

FastAPI OpenAI Chat Completion接口问题解决方案

问题1:流式会话被非流式请求中断

核心原因

流式接口若使用同步阻塞逻辑或共享全局非线程安全资源,会占用FastAPI的事件循环,导致新的非流式请求抢占资源,中断正在进行的流式传输。

解决方案

  1. 强制异步处理流式任务
    所有流式生成逻辑必须异步化,避免阻塞事件循环。若调用的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")
    
  2. 隔离请求级资源
    禁止在不同请求间共享LLM客户端实例、会话状态等非线程安全对象,通过依赖注入为每个请求创建独立实例。


问题2:流式API无请求完成日志且连接未断开

核心原因

流式响应默认在生成器耗尽后关闭连接,但缺少明确的结束标记和日志回调,导致客户端可能维持连接,且FastAPI未打印请求完成日志。

解决方案

  1. 添加流式结束标记
    在生成器末尾输出标准结束信号,告知客户端会话完成:

    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"
    
  2. 自定义中间件打印完成日志
    通过中间件包装流式响应的生成器,在生成器耗尽后打印与非流式一致的请求日志:

    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 19:44:54