Python + FastAPI 异步流式对话接口设计
共 8610字,需浏览 18分钟
·
2026-07-02 22:39
在大模型应用里,普通的同步 HTTP 接口很快会遇到体验和工程上的双重瓶颈:用户要等完整回答生成完才能看到结果,服务端也很难在一次长请求中清晰处理超时、限流、日志追踪、会话状态和 token 统计。
更稳妥的方式,是把对话生成设计成一个异步流式接口。客户端发起一次请求,服务端立即返回 text/event-stream,随后持续推送模型生成过程中的事件。这样既能改善首字延迟,也能把请求追踪、会话互斥、思考模式、超时和 usage 统计纳入统一的接口契约。
本文设计一个基于 Python + FastAPI 的单接口方案。
目标
这个接口需要解决以下问题:
- 支持异步流式输出。
- 支持思考模式,但不直接暴露原始隐藏推理链。
- 支持限流,避免用户、IP 或 session 维度的并发滥用。
- 支持结构化日志,能通过
request_id串起完整链路。 - 支持超时控制,避免模型调用长时间占用连接。
- 获取或生成
request_id。 - 获取或创建对应的
session_id。 - 判断同一 session 中是否还有未结束的对话。
- 统计输入、输出和总 token 数。
接口定义
推荐只设计一个对外接口:
POST /v1/chat/stream
Content-Type: application/json
Accept: text/event-stream
X-Request-ID: 可选
Authorization: Bearer <token>请求体示例:
{
"session_id": "sess_01JZ2W6G4A8V0Y9H6D7M2KQ5PV",
"message": "帮我分析这个面试案例",
"thinking": {
"enabled": true,
"budget_tokens": 1024,
"visible_summary": true
},
"stream_options": {
"timeout_seconds": 60,
"include_usage": true
}
}其中:
session_id可选。客户端不传时,服务端创建新 session,并在流式事件中返回。message是用户本轮输入。thinking.enabled表示启用更强推理模式。thinking.budget_tokens表示允许模型用于推理的预算。thinking.visible_summary表示是否向前端发送“思考摘要”或“思考状态”事件。stream_options.timeout_seconds是本次生成的业务超时时间,服务端应设置最大上限。stream_options.include_usage控制最后是否返回 token usage。
响应类型:
HTTP/1.1 200 OK
Content-Type: text/event-stream
Cache-Control: no-cache
X-Request-ID: req_01JZ2W6JQX89S6JDQ6D8C0K3AS流式事件设计
建议统一使用 Server-Sent Events,也就是 SSE。每个事件包含 event 和 data:
event: meta
data: {"request_id":"req_xxx","session_id":"sess_xxx"}
event: thinking
data: {"type":"status","text":"正在分析问题结构"}
event: delta
data: {"content":"这是"}
event: delta
data: {"content":"一个"}
event: usage
data: {"input_tokens":128,"output_tokens":642,"total_tokens":770}
event: done
data: {"finish_reason":"stop"}推荐事件类型如下:
| event | 作用 |
|---|---|
meta | 返回 request_id、session_id、限流信息等元数据。 |
thinking | 返回思考状态或思考摘要,不返回原始隐藏推理链。 |
delta | 返回模型正文增量。 |
usage | 返回 token 统计。 |
error | 返回业务错误。 |
done | 表示本次流结束。 |
这里要特别注意思考模式。接口可以支持“思考中”“已完成分析”“摘要如下”这类可见状态,但不建议把模型的原始 chain-of-thought 逐字返回给前端。更合适的做法是把思考模式设计成模型参数和产品状态,而不是把内部推理过程当作普通内容暴露。
request_id 设计
request_id 是日志追踪的核心。推荐规则:
- 优先读取客户端请求头
X-Request-ID。 - 如果没有传,服务端生成一个新的 ID。
- 校验 ID 长度和字符集,防止日志注入。
- 响应头和首个
meta事件都返回同一个request_id。 - 后续所有日志都带上这个
request_id。
示例伪代码:
def get_request_id(request: Request) -> str:
request_id = request.headers.get("X-Request-ID")
if request_id and is_safe_request_id(request_id):
return request_id
return generate_id("req")这样客户端、网关、服务端日志和模型调用日志就能串成一条完整链路。
session_id 与会话状态
session_id 用于标识一段连续对话。客户端可以传入已有 session,也可以不传,让服务端创建。
设计上需要维护 session 的状态:
session_id
user_id
status: idle | generating
active_request_id
created_at
updated_at
expires_at在请求开始时,服务端需要判断同一个 session 是否已有未结束的生成任务:
- 如果 session 不存在,创建 session。
- 如果 session 存在且状态为
idle,把它更新为generating。 - 如果 session 存在且状态为
generating,说明上一轮对话还没结束,返回409 Conflict。 - 流式生成结束、超时、客户端断开或异常时,都必须把 session 状态恢复为
idle。
这个判断不能只放在内存里。生产环境建议使用 Redis 分布式锁:
lock key: chat:session:{session_id}:active
value: request_id
ttl: timeout_seconds + buffer_seconds获取锁成功,说明可以开始生成;获取锁失败,说明这个 session 还有对话没结束。
限流设计
限流建议至少包含三层:
| 维度 | 目的 |
|---|---|
user_id | 控制单用户总体请求量。 |
ip | 防止匿名或异常来源刷接口。 |
session_id | 防止同一会话并发生成。 |
对流式接口来说,限流不仅要限制 QPS,还要限制并发连接数。一个用户如果同时打开很多长连接,会比普通短请求更消耗资源。
推荐策略:
user_id:每分钟请求数限制,例如30/min。ip:每分钟请求数限制,例如60/min。user_id active streams:并发流限制,例如最多3。session_id:同一时间只允许1个 active stream。
如果触发限流,可以直接返回:
HTTP/1.1 429 Too Many Requests
Retry-After: 10响应体:
{
"error": {
"code": "rate_limited",
"message": "请求过于频繁,请稍后再试",
"request_id": "req_xxx"
}
}如果请求已经进入 SSE 流之后才发生下游限流,则通过 error 事件返回,再发送 done 事件结束流。
超时设计
流式接口至少需要三类超时:
| 超时 | 说明 |
|---|---|
| request timeout | 整个接口允许执行的最长时间。 |
| upstream timeout | 调用模型服务的最长等待时间。 |
| idle timeout | 多久没有任何 token 或事件返回就中断。 |
设计建议:
- 客户端可以传
timeout_seconds,但服务端必须设置最大值,例如120秒。 - 服务端通过
asyncio.timeout()包裹主生成逻辑。 - 对模型流式返回增加 idle timer,避免连接挂死。
- 超时后发送
error事件,错误码为timeout。 - 无论正常结束还是超时,都执行清理逻辑,释放 session 锁和用户并发计数。
token 计数
token 统计分为三部分:
{
"input_tokens": 128,
"reasoning_tokens": 256,
"output_tokens": 642,
"total_tokens": 1026
}设计建议:
input_tokens在调用模型前计算,包含系统提示词、历史消息和当前用户输入。output_tokens可以在流式输出过程中累加,也可以在结束后统一计算。reasoning_tokens如果模型供应商返回该字段,则使用供应商返回值;如果没有,不要伪造精确值。total_tokens = input_tokens + reasoning_tokens + output_tokens,但要注意不同模型供应商的 usage 字段可能已经包含 reasoning tokens。
如果业务上要做计费,建议以模型供应商最终返回的 usage 为准;本地 token 计数只用于预估、限额和日志。
日志设计
流式接口的日志要记录关键节点,而不是把所有 token 都写进日志。
推荐日志字段:
{
"request_id": "req_xxx",
"session_id": "sess_xxx",
"user_id": "user_xxx",
"event": "chat_stream_finished",
"latency_ms": 12840,
"input_tokens": 128,
"output_tokens": 642,
"finish_reason": "stop",
"error_code": null
}推荐记录这些事件:
chat_stream_startedsession_lock_acquiredrate_limit_checkedmodel_stream_startedfirst_token_receivedchat_stream_finishedchat_stream_failedsession_lock_released
不要默认记录完整用户输入和完整模型输出。如果需要排查问题,可以记录摘要、长度、hash 或经过脱敏后的内容。
FastAPI 接口骨架
下面是一个偏设计性质的骨架,展示接口边界和控制流:
from fastapi import APIRouter, HTTPException, Request
from fastapi.responses import StreamingResponse
router = APIRouter()
@router.post("/v1/chat/stream")
async def chat_stream(request: Request, payload: ChatStreamRequest):
request_id = get_request_id(request)
user_id = get_current_user_id(request)
session_id = payload.session_id or create_session_id()
await check_rate_limit(
user_id=user_id,
session_id=session_id,
request_id=request_id,
)
lock_acquired = await acquire_session_lock(
session_id=session_id,
request_id=request_id,
ttl_seconds=payload.effective_timeout + 10,
)
if not lock_acquired:
raise HTTPException(
status_code=409,
detail={
"code": "session_busy",
"message": "当前会话中还有未结束的对话",
"request_id": request_id,
"session_id": session_id,
},
)
async def event_generator():
usage = TokenUsage()
try:
yield sse("meta", {
"request_id": request_id,
"session_id": session_id,
})
async with timeout(payload.effective_timeout):
usage.input_tokens = count_input_tokens(payload)
async for chunk in model_stream(
payload=payload,
request_id=request_id,
session_id=session_id,
):
if chunk.type == "thinking":
yield sse("thinking", to_visible_thinking_event(chunk))
elif chunk.type == "content":
usage.output_tokens += count_output_tokens(chunk.text)
yield sse("delta", {"content": chunk.text})
elif chunk.type == "usage":
usage.merge(chunk.usage)
yield sse("usage", usage.to_dict())
yield sse("done", {"finish_reason": "stop"})
except TimeoutError:
yield sse("error", {
"code": "timeout",
"message": "生成超时",
"request_id": request_id,
})
yield sse("done", {"finish_reason": "timeout"})
except Exception:
log_exception(request_id=request_id, session_id=session_id)
yield sse("error", {
"code": "internal_error",
"message": "生成失败",
"request_id": request_id,
})
yield sse("done", {"finish_reason": "error"})
finally:
await release_session_lock(
session_id=session_id,
request_id=request_id,
)
await decrease_active_stream_count(user_id)
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={"X-Request-ID": request_id},
)这个骨架的重点不是模型调用细节,而是接口的生命周期:
- 获取
request_id。 - 获取或创建
session_id。 - 做限流检查。
- 判断 session 是否已有未结束对话。
- 返回 SSE 流。
- 生成过程中推送
thinking、delta、usage等事件。 - 正常结束、异常或超时后释放锁。
错误码设计
建议提前定义稳定错误码,方便前端和调用方处理:
| code | HTTP 状态 | 场景 |
|---|---|---|
invalid_request | 400 | 参数错误。 |
unauthorized | 401 | 未认证。 |
rate_limited | 429 | 触发限流。 |
session_busy | 409 | 同一 session 还有生成未结束。 |
timeout | 200 SSE / 504 | 生成超时。 |
upstream_error | 200 SSE / 502 | 模型服务异常。 |
internal_error | 200 SSE / 500 | 服务内部错误。 |
对于流式接口,有一个设计细节:如果错误发生在响应头发送之前,可以直接返回 HTTP 错误码;如果错误发生在 SSE 已经开始之后,就只能通过 event: error 告知客户端,然后用 event: done 结束流。
客户端处理建议
客户端需要按事件类型分别处理:
- 收到
meta后保存request_id和session_id。 - 收到
thinking后更新“思考中”状态或展示思考摘要。 - 收到
delta后追加正文。 - 收到
usage后展示或上报 token 使用量。 - 收到
error后展示错误状态。 - 收到
done后关闭连接并允许用户继续下一轮提问。
如果客户端断开连接,服务端也要感知并清理状态。FastAPI 中可以通过 await request.is_disconnected() 在生成循环里检查连接是否仍然存在。
总结
一个可靠的异步流式对话接口,不只是把模型输出改成 streaming。真正重要的是把请求生命周期设计完整:request_id 负责链路追踪,session_id 负责对话归属,session lock 负责防止同一会话并发生成,限流负责保护系统资源,超时负责释放长连接,日志负责问题排查,token usage 负责成本统计和计费基础。
如果只做一个接口,推荐把它收敛成:
POST /v1/chat/stream再通过稳定的 SSE 事件协议承载 meta、thinking、delta、usage、error 和 done。这样接口表面保持简单,内部却能覆盖真实生产环境最需要的控制点。
