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

Google Cloud Run上FastAPI流式传输OpenAI响应至Flutter应用时出现504超时错误

问题:FastAPI + Cloud Run 流式请求频繁504超时及线程阻塞排查

我将一个仅包含三个Python文件的FastAPI项目部署到Google Cloud Run,通过SSE把OpenAI的响应流式传输给Flutter移动端。Cloud Run的请求超时配置为31秒,移动端自身设置10秒超时——若10秒内未收到流数据则判定请求失败。当前服务器出现异常:有时所有请求都会返回504超时,持续约5分钟,之后会出现一系列Traceback日志。怀疑是线程阻塞导致请求排队,但无法定位问题,恳请帮忙审查代码并提供排查思路。


项目代码

main.py

from dependencies import denpendency
from routers import chat_v1
import logging
from fastapi import FastAPI, Request, status, Depends
from fastapi.exceptions import RequestValidationError
from fastapi.responses import JSONResponse
from starlette.exceptions import HTTPException as StarletteHTTPException
from fastapi.responses import PlainTextResponse

app = FastAPI(
    title='FastAPI API',
    openapi_url='/',
    docs_url='/docs',
    dependencies=[Depends(denpendency.validate_jwt)],
)
app.include_router(chat_v1.router)


@app.exception_handler(StarletteHTTPException)
async def http_exception_handler(request, exc):
    return PlainTextResponse(str(exc.detail), status_code=exc.status_code)


@app.exception_handler(RuntimeError)
async def runtime_exception_handler(request, exc):
    return PlainTextResponse(str(exc.detail), status_code=exc.status_code)


@app.exception_handler(Exception)
async def exception_handler(request, exc):
    return PlainTextResponse(str(exc))


@app.exception_handler(RequestValidationError)
async def validation_exception_handler(request: Request, exc: RequestValidationError):
    exc_str = f'{exc}'.replace('\n', ' ').replace('   ', ' ')
    logging.error(f"{request}: {exc_str}")
    content = {'status_code': 10422, 'message': exc_str, 'data': None}
    return JSONResponse(content=content, status_code=status.HTTP_422_UNPROCESSABLE_ENTITY)

denpendency.py

from fastapi import Header, HTTPException, Depends
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer
from firebase_admin import auth, credentials, initialize_app

cred = credentials.Certificate("assets/credentials/sa.json")
default_app = initialize_app(cred)
print("Name: " + default_app.name)


def validate_jwt(credential: HTTPAuthorizationCredentials = Depends(HTTPBearer(auto_error=False))):
    if cred is None:
        raise HTTPException(
            status_code=401,
            detail="Bearer authentication is needed",
            headers={'WWW-Authenticate': 'Bearer realm="auth_required"'},
        )
    try:
        decoded_token = auth.verify_id_token(credential.credentials)
    except Exception as err:
        raise HTTPException(
            status_code=401,
            detail=f"Invalid authentication from Firebase. {err}",
            headers={'WWW-Authenticate': 'Bearer error="invalid_token"'},
        )
    return decoded_token

chat_v1.py

async def get_answer_stream(request: Request, body: ChatRequest):
    try:
        message_stream = openai.ChatCompletion.create(
            model='gpt-3.5-turbo-16k',
            messages=get_messages(body),
            stream=True,
        )
    except Exception as e:
        error_message = "OpenAI Response (Streaming) Error: " + str(e)
        print(error_message)
        return

    try:
        for chunk in message_stream:
            if await request.is_disconnected():
                print("Disconnected!")
                break
            choices = chunk["choices"][0]
            delta = choices["delta"]
            if 'content' in delta:
                content = chunk["choices"][0]["delta"]['content']
                done = False
            else:
                content = None
                done = 'finish_reason' in choices and choices['finish_reason'] == 'stop'
            yield json.dumps({"token": content, "done": done})
    except CancelledError:
        print("User stop generating message")
    except Exception as e:
        print("OpenAI Response (Streaming) Error: " + str(e))
    finally:
        print('Finish Stream flow')
        await request.close()


@router.post(
    "/getAnswer",
    tags=["chat-v1"],
    response_model=str,
    responses={503: {"detail": error503}}
)
async def get_answer_v1(request: Request):
    body = await request.json()
    chat_request = ChatRequest(
        messages=body.get('messages'),
        tag=body.get('tag'),
        prompts=body.get('prompts'),
    )
    answer = get_answer_stream(request, chat_request)
    # from sse_starlette.sse import EventSourceResponse
    result = EventSourceResponse(answer)
    return result

错误日志

504超时日志

{
//504 error message
httpRequest: {10}
insertId: "64ac6a5c000aabf6c037a3b0"
labels: {1}
logName: "projects/projectname/logs/run.googleapis.com%2Frequests"
receiveTimestamp: "2023-07-10T20:30:20.846742621Z"
resource: {2}
severity: "ERROR"
spanId: "9149462740768920648"
textPayload: "The request has been terminated because it has reached the maximum request timeout. To change this limit, see https://cloud.google.com/run/docs/configuring/request-timeout"
timestamp: "2023-07-10T20:29:49.696909Z"
trace: "projects/projectname/traces/49407a47b7bb6d0953078dbf9416aa1b"
}

Traceback日志

Traceback (most recent call last): 
  File "/usr/local/lib/python3.9/site-packages/uvicorn/protocols/http/h11_impl.py", line 373, in run_asgi
    result = await app(self.scope, self.receive, self.send) 
  File "/usr/local/lib/python3.9/site-packages/uvicorn/middleware/proxy_headers.py", line 75, in __call__
    return await self.app(scope, receive, send) 
  File "/usr/local/lib/python3.9/site-packages/fastapi/applications.py", line 208, in __call__
    await super().__call__(scope, receive, send) 
  File "/usr/local/lib/python3.9/site-packages/starlette/applications.py", line 112, in __call__
    await self.middleware_stack(scope, receive, send) 
  File "/usr/local/lib/python3.9/site-packages/starlette/middleware/errors.py", line 181, in __call__
    raise exc from None 
  File "/usr/local/lib/python3.9/site-packages/starlette/middleware/errors.py", line 159, in __call__
    await self.app(scope, receive, _send) 
  File "/usr/local/lib/python3.9/site-packages/starlette/exceptions.py", line 82, in __call__
    raise exc from None 
  File "/usr/local/lib/python3.9/site-packages/starlette/exceptions.py", line 71, in __call__
    await self.app(scope, receive, sender) 
  File "/usr/local/lib/python3.9/site-packages/starlette/routing.py", line 580, in __call__
    await route.handle(scope, receive, send) 
  File "/usr/local/lib/python3.9/site-packages/starlette/routing.py", line 241, in handle
    await self.app(scope, receive, send) 
  File "/usr/local/lib/python3.9/site-packages/starlette/routing.py", line 52, in app
    response = await func(request) 
  File "/usr/local/lib/python3.9/site-packages/fastapi/routing.py", line 226, in app
    raw_response = await run_endpoint_function( 
  File "/usr/local/lib/python3.9/site-packages/fastapi/routing.py", line 159, in run_endpoint_function
    return await dependant.call(**values) 
  File "/app/./routers/chat_v1.py", line 106, in get_answer_v1
    body = await request.json() 
  File "/usr/local/lib/python3.9/site-packages/starlette/requests.py", line 227, in json
    body = await self.body() 
  File "/usr/local/lib/python3.9/site-packages/starlette/requests.py", line 220, in body
    async for chunk in self.stream(): 
  File "/usr/local/lib/python3.9/site-packages/starlette/requests.py", line 214, in stream
    raise ClientDisconnect()

错误日志截图

错误日志截图


排查思路与修复方案

一、核心问题定位

从Traceback的ClientDisconnect错误结合504超时来看,核心问题集中在:

  1. 同步OpenAI调用阻塞事件循环,导致请求排队
  2. 客户端断开后资源未及时释放
  3. JWT验证逻辑存在变量名错误
  4. 手动读取请求体引发的异常阻塞

二、代码修复点

1. 替换同步OpenAI调用为异步版本

当前openai.ChatCompletion.create是同步方法,会直接阻塞FastAPI的事件循环,导致后续请求无法处理。改用异步OpenAI客户端:

# 先安装异步依赖:pip install openai[async]
from openai import AsyncOpenAI

client = AsyncOpenAI()

async def get_answer_stream(request: Request, body: ChatRequest):
    try:
        # 异步调用OpenAI流式接口
        message_stream = await client.chat.completions.create(
            model='gpt-3.5-turbo-16k',
            messages=get_messages(body),
            stream=True,
        )
    except Exception as e:
        error_message = "OpenAI Response (Streaming) Error: " + str(e)
        print(error_message)
        yield json.dumps({"token": None, "done": True, "error": error_message})
        return

    try:
        # 异步遍历流数据
        async for chunk in message_stream:
            if await request.is_disconnected():
                print("Disconnected!")
                await message_stream.aclose()
                break
            choices = chunk.choices[0]
            delta = choices.delta
            if delta.content is not None:
                content = delta.content
                done = False
            else:
                content = None
                done = choices.finish_reason == 'stop'
            yield json.dumps({"token": content, "done": done})
    except CancelledError:
        print("User stop generating message")
        await message_stream.aclose()
    except Exception as e:
        print("OpenAI Response (Streaming) Error: " + str(e))
    finally:
        print('Finish Stream flow')

2. 修复JWT验证的变量名错误

denpendency.py中存在变量名混淆:判断的是全局变量cred而非请求参数credential,导致逻辑错误:

def validate_jwt(credential: HTTPAuthorizationCredentials = Depends(HTTPBearer(auto_error=False))):
    # 原错误:if cred is None:
    if credential is None:
        raise HTTPException(
            status_code=401,
            detail="Bearer authentication is needed",
            headers={'WWW-Authenticate': 'Bearer realm="auth_required"'},
        )
    # 延迟初始化Firebase,避免模块导入时的冷启动问题
    if not firebase_admin._apps:
        cred = credentials.Certificate("assets/credentials/sa.json")
        initialize_app(cred)
    try:
        decoded_token = auth.verify_id_token(credential.credentials)
    except Exception as err:
        raise HTTPException(
            status_code=401,
            detail=f"Invalid authentication from Firebase. {err}",
            headers={'WWW-Authenticate': 'Bearer error="invalid_token"'},
        )
    return decoded_token

3. 改用依赖注入自动解析请求体

手动读取request.json()容易在客户端断开时引发异常,改用FastAPI的依赖注入自动解析ChatRequest:

@router.post(
    "/getAnswer",
    tags=["chat-v1"],
    responses={503: {"detail": error503}}
)
async def get_answer_v1(request: Request, chat_request: ChatRequest = Depends()):
    # 直接使用注入的chat_request,无需手动解析
    answer = get_answer_stream(request, chat_request)
    result = EventSourceResponse(answer)
    return result

三、Cloud Run配置优化

  • 调整并发数:默认并发数80,可降至30-50,避免单个实例承载过多流式请求导致阻塞
  • 冷启动优化:设置最小实例数为1,避免冷启动时的请求排队;同时压缩镜像大小,减少启动时间
  • 超时配置:若OpenAI流式响应经常超过31秒,可适当延长Cloud Run请求超时(最大900秒),同时移动端需配合实现重试逻辑

四、排查工具

  • 查看Cloud Run监控面板:关注请求队列长度、实例数变化、CPU/内存使用率,判断是否为资源不足导致阻塞
  • 添加详细日志:在请求开始/结束、OpenAI调用前后记录时间戳,定位慢请求
  • 事件循环调试:测试环境下启动时添加uvicorn.run(app, debug=True),查看事件循环阻塞情况

内容的提问来源于stack exchange,提问作者Nam Đỗ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 05:34:52