在 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. 每个任务拥有独立事件队列

miniagentProgressTracker 使用:

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 架构

flowchart LR A["Document / OCR / KB Task"] -->|"emit()"| B["ProgressTracker"] B --> C["asyncio.Queue
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实现数据流图:

flowchart TD U["Browser / Frontend"] -->|"1. Start Task"| API["FastAPI Task API"] API -->|"2. return task_id"| U API -->|"3. create task"| TASK["Document / OCR / KB Task"] TASK -->|"4. emit()"| PT["ProgressTracker"] PT --> Q["asyncio.Queue
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 编程,我们要避免的不只是错误代码,还包括这种架构浪费:

明明服务器可以主动告诉你结果,却让客户端不停地问:“好了吗?”

这正是长连接与流式推送存在的价值。


开源代码


🪐祝您好运🪐