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超时来看,核心问题集中在:
- 同步OpenAI调用阻塞事件循环,导致请求排队
- 客户端断开后资源未及时释放
- JWT验证逻辑存在变量名错误
- 手动读取请求体引发的异常阻塞
二、代码修复点
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 Đỗ
相关产品推荐
相关产品推荐

