在 Web 系统里,有一类需求特别容易被 AI 写成“能跑,但很笨”的实现:
前端:
每 1 秒问一次服务器
“任务完成了吗?”
服务器:
“没有。”
1 秒后再问:
“完成了吗?”
这种方式叫:
Polling:轮询
偶尔查询一次状态没有问题,但如果是 AI 流式回答、OCR 进度、文档解析、知识库构建、后台任务状态等实时场景,高频轮询会产生大量无效请求。
更合适的方案通常是:
SSE(Server-Sent Events,服务器发送事件)
或者:
WebSocket:Web 套接字,全双工长连接通信协议
核心思想只有一句:
不要让客户端不停问“有消息了吗”,而是让服务器有消息时主动推送。
📌 技术名片
Long-lived Connection:长连接
客户端和服务器建立连接后,不立即关闭,而是在一段时间内持续保持连接,用于后续数据传输。
常见方案主要有两种。
SSE
SSE(Server-Sent Events,服务器发送事件) 基于 HTTP,主要用于:
Server
↓
Client
适合:
- AI Token 流式输出;
- 后台任务进度;
- 日志流;
- 状态通知。
WebSocket
WebSocket:Web 套接字协议 支持:
Client
↕
Server
也就是:
Full-Duplex Communication:全双工通信
更适合:
- 在线聊天;
- 多人协作;
- 实时游戏;
- 双向实时控制。
简单判断:
只需要服务器推给客户端,优先 SSE;双方都要高频主动发消息,再考虑 WebSocket。
💡 一个简单比喻
轮询像不停给维修店打电话:
你:修好了吗?
店员:没有。
1 分钟后:
你:修好了吗?
店员:还没有。
而服务器推送更像:
你:修好了通知我。
……
店员:修好了。
所以可以记住:
Polling
客户端不断问
SSE
服务器有消息就推送
WebSocket
双方随时都能讲话
一、AI 最容易犯的错误:用高频轮询模拟实时
假设让 AI 实现:
上传 PDF 后实时显示处理进度。
它很容易生成:
POST /tasks
GET /tasks/{task_id}/status
然后前端:
setInterval(async () => {
const result = await getTaskStatus(taskId)
}, 1000)
如果 100 个用户每秒查询一次,服务器每秒就多出:
100 HTTP Requests
一个任务持续 60 秒:
100 × 60 = 6000 次请求
但真正发生状态变化的次数可能只有几十次。
因此:
不要用高频轮询模拟实时推送。
二、架构规范
1. 任务和流式传输要解耦
不推荐:
DocumentService
↓
直接控制 StreamingResponse
↓
直接 yield SSE
更合理的是:
Business Task
↓
Progress Event
↓
Async Queue
↓
SSE Endpoint
↓
Browser
也就是说:
任务负责产生事件,SSE 接口负责发送事件。
这样业务逻辑不需要知道 HTTP 协议细节。
2. 统一事件格式
不同任务不要各自定义:
{"percent": 30}
或者:
{"status": "almost_done"}
推荐统一成:
{
"stage": "embed",
"message": "Embedding chunks...",
"progress": 72.5,
"done": false,
"error": false,
"ts": "2026-08-17T10:00:00"
}
含义:
stage
当前阶段
message
状态说明
progress
进度 0~100
done
是否完成
error
是否失败
ts
事件时间
这样前端可以复用同一套进度组件。
3. SSE 生命周期必须完整
一个可靠的 SSE 实现至少要考虑四件事:
创建 Queue
↓
持续发送 Event
↓
Keepalive 保活
↓
done / error 后关闭并清理
Keepalive:保活
长时间没有数据时,中间代理可能误以为连接失效。
因此可以周期性发送:
: keepalive
明确结束条件
任务完成:
done = true
任务失败:
error = true
都应该结束流。
清理资源
连接关闭后要移除任务队列,否则长期运行可能造成内存积累。
队列必须有上限
例如:
asyncio.Queue(maxsize=200)
防止客户端消费过慢时事件无限堆积。
这与:
Backpressure:背压
机制有关,即消费者跟不上生产者时,系统要有办法限制积压。
4. 关闭缓存和代理缓冲
SSE 常见问题之一是:
服务器已经持续:
yield event
浏览器却迟迟收不到,最后一次性收到很多条。
通常是代理缓冲造成的,因此常见响应配置是:
Cache-Control: no-cache
X-Accel-Buffering: no
目的很简单:
事件产生后尽快送到客户端,而不是积攒后再发送。
5. 任务执行和任务观察要分离
对于 OCR、文档处理、Embedding 等长任务,更合理的结构是:
Client
↓
Start Task
↓
task_id
Background Task
↓
持续产生 Progress Event
Client
↓
SSE /tasks/{task_id}/progress
即:
任务执行负责做事,SSE 负责观察进度。
两者通过 task_id 关联。
6. 为什么不要滥用 WebSocket
WebSocket 功能更强,但也意味着需要额外处理:
- 心跳;
- 重连;
- 身份认证;
- 连接管理;
- 多实例路由;
- 消息顺序;
- 广播;
- 断线清理。
如果业务只是:
20%
40%
80%
100%
这种单向进度推送,SSE 往往更简单。
所以:
架构不是 WebSocket 越多越先进,而是协议与通信模式匹配。
三、长连接与流式推送有什么好处?
1. 减少无效请求
轮询:
Client → Server
Client → Server
Client → Server
SSE:
Client ───────── Server
↓
20%
↓
50%
↓
100%
只有真正有变化时才推送业务事件。
2. 降低感知延迟
Perceived Latency:感知延迟
AI 一次回答即使总耗时仍是 10 秒:
传统模式:
等待 10 秒
↓
完整答案突然出现
流式模式:
第 1 秒开始看到内容
↓
持续输出
↓
第 10 秒结束
用户会明显感觉系统更快。
3. 非常适合 AI 系统
AI 系统天然包含:
LLM Token Stream
文档处理进度
Embedding 进度
Web Search 状态
工具执行过程
Agent 执行状态
所以流式推送通常是 Agent 系统的重要基础能力。
四、提示词落地:直接约束 AI
可以把以下规则写进:
Project Rules:项目级规则
## 流式传输与长连接规则
1. 当服务器推送更合适时,请勿使用高频轮询来获取实时进度。
2. 优先使用 SSE 进行服务器到客户端的流式传输:
- AI 令牌流式传输
- 后台任务进度
- OCR/PDF 处理进度
- 日志
- 通知
3. 仅在需要频繁双向通信时才使用 WebSocket
4. 将任务执行与流式传输分离。
推荐架构:任务 → 进度事件 → 异步队列 → SSE 端点 → 客户端
5. 使用统一的事件模式:
{
stage,
message,
progress,
done,
error,
ts
}
6. SSE 必须支持 keepalive。
7. 当 done=true 或 error=true 时关闭流。
8. 终止后清理队列和资源。
9. 事件队列必须具有有限的容量。
10. 必要时禁用 SSE 的响应缓冲。
11. 在生成轮询代码之前,请检查 SSE或 WebSocket 哪种方式更合适。
这比一句:
“帮我实现实时进度。”
更能避免 AI 下意识生成 setInterval()。
五、正面产出:miniagent 的真实 SSE 实现
miniagent 已经实现了一套完整的任务进度 SSE 架构:
Task
↓
ProgressTracker
↓
asyncio.Queue
↓
SSE Endpoint
↓
Browser
1. 每个任务拥有独立事件队列
miniagent 的 ProgressTracker 使用:
class ProgressTracker:
_queues: dict[str, asyncio.Queue] = {}
@classmethod
def create(cls, task_id: str) -> asyncio.Queue:
q: asyncio.Queue = asyncio.Queue(maxsize=200)
cls._queues[task_id] = q
return q
每个任务对应自己的 Queue:
task_001 → Queue A
task_002 → Queue B
task_003 → Queue C
并通过:
maxsize=200
限制事件积压。
2. 业务层只负责发布事件
miniagent 通过:
await ProgressTracker.emit(
task_id,
stage="embed",
message="Embedding chunks...",
progress=72.5,
)
发布统一结构:
{
"stage": stage,
"message": message,
"progress": round(progress, 1),
"done": done,
"error": error,
"ts": datetime.now().isoformat(),
}
任务本身不需要操作 StreamingResponse,实现了业务逻辑与 SSE 传输层解耦。
3. SSE Endpoint 持续消费事件
miniagent 提供:
GET /{task_id}/progress
核心逻辑:
queue = ProgressTracker.get(task_id)
event = await asyncio.wait_for(
queue.get(),
timeout=30.0,
)
有新进度时,Queue 立即把事件交给 SSE Endpoint。
4. 30 秒无事件时发送 Keepalive
miniagent 在超时时:
except asyncio.TimeoutError:
yield ": keepalive\n\n"
让长连接在任务暂时没有新状态时仍保持活跃。
5. 完成或失败后自动关闭并清理
发送事件:
yield f"data: {json.dumps(event, ensure_ascii=False)}\n\n"
然后判断:
if event.get("done") or event.get("error"):
break
最后:
finally:
ProgressTracker.remove(task_id)
也就是说完整生命周期是:
Create Queue
↓
Emit Progress
↓
SSE Push
↓
done / error
↓
Close Stream
↓
Remove Queue
6. 禁止缓存与代理缓冲
最终 miniagent 返回:
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
},
)
其中:
text/event-stream
表示 SSE 数据流。
而:
Cache-Control: no-cache
X-Accel-Buffering: no
用于减少缓存和反向代理缓冲对实时性的影响。
六、miniagent 的 SSE 架构
maxsize = 200"] D["Browser"] -->|"GET /tasks/{id}/progress"| E["FastAPI SSE Endpoint"] C -->|"await queue.get()"| E E -->|"data: JSON"| D E -. "30s timeout" .-> F["keepalive"] E -->|"done / error"| G["Close Stream"] G --> H["Remove Queue"]
这套结构最重要的边界是:
任务生产者
↓
事件队列
↓
SSE 传输层
↓
客户端
而不是让任务代码直接控制 HTTP Response。
下面是更详细的SSE实现数据流图:
task_id → Queue
maxsize=200"] U -->|"5. GET /{task_id}/progress"| SSE["FastAPI SSE Endpoint"] Q -->|"6. await queue.get()"| SSE SSE -->|"7. data: JSON"| U SSE -. "30s timeout" .-> KA["Keepalive"] KA -.-> U TASK -->|"progress / stage / message"| PT TASK -->|"done / error"| PT SSE -->|"8. done / error"| CLOSE["Close Stream"] CLOSE --> CLEAN["9. Remove Queue"]
七、给 AI 的检查清单
以后让 AI 实现实时功能时,可以要求它检查:
Streaming Architecture Checklist
□ 这里真的需要轮询吗?
□ 如果只是 Server → Client,是否优先使用 SSE?
□ 是否真的需要双向 WebSocket?
□ 任务执行是否与 Streaming 解耦?
□ 是否使用统一事件格式?
□ Queue 是否有容量限制?
□ SSE 是否有 keepalive?
□ done / error 后是否结束?
□ 结束后是否清理资源?
□ 是否禁用了缓存和代理缓冲?
总结
AI 很容易生成:
setInterval
↓
GET /status
↓
没完成
↓
再请求
对于实时进度、AI 输出和长任务,这往往会产生大量无效请求。
更加合理的设计是:
后台任务
↓
Progress Event
↓
Async Queue
↓
SSE
↓
Browser
只有真正需要频繁双向通信时,再使用:
Client
↕
WebSocket
↕
Server
miniagent 当前已经通过 ProgressTracker、有界 asyncio.Queue、FastAPI StreamingResponse、30 秒 Keepalive、done/error 终止以及 Queue 清理,形成了一套完整的 SSE 任务进度推送链路。
对于 AI 编程,我们要避免的不只是错误代码,还包括这种架构浪费:
明明服务器可以主动告诉你结果,却让客户端不停地问:“好了吗?”
这正是长连接与流式推送存在的价值。
开源代码
🪐祝您好运🪐