返回文章列表
frontend2026年9月21日约 16 分钟阅读

SSE 业务建模:从协议到生产中的前后端核心链路

以 AI 写作任务为例,从 SSE 的消息格式出发,串起创建 Job、Worker 执行、事件持久化、服务端推送和前端消费,再理解断线重连、去重与生产边界。

用户点击“生成文章”,页面开始展示进度,正文逐段出现。中途断网、刷新页面,再回来还能继续看到结果。这套体验背后,SSE 负责哪一段?任务是谁执行的,内容存在哪里,前端又如何恢复?

我想把这条链路从头梳理清楚:先理解 SSE 怎样传消息,再给任务和事件建模,最后把后端产生事件到前端更新页面的代码连起来。

全文沿用一个 AI 写作任务。代码是围绕生产问题提炼的 TypeScript 核心骨架,数据库、队列和 HTTP 框架通过接口表示,不是一份复制后即可部署的完整工程。

1. SSE 是什么:一个持续返回事件的 HTTP 响应

普通接口通常在计算结束后返回一次 JSON。SSE(Server-Sent Events)让服务端保持 HTTP 响应打开,持续往响应体里写事件,浏览器边接收边处理。

浏览器                                  服务端
   ─── GET /jobs/job_001/events ──────────→
   ←── 200,Content-Type: text/event-stream
   ←── progress:正在分析资料
   ←── content:第一段正文
   ←── content:第二段正文
   ←── done:生成完成

它适合任务进度、生成内容、通知这类服务端向客户端推送的场景。客户端提交输入、取消任务,仍然可以走普通 HTTP 接口;SSE 这条响应只承担向下发送事件的方向。

一次响应中可以包含许多事件。下面是一条实际的线格式,末尾必须有空行:

id: 3
event: content
data: {"text":"第一段正文"}
 

id 标识事件,event 表示类型,data 是文本,这里约定使用 JSON。SSE 使用 UTF-8 编码;JSON、任务状态和内容追加规则都是应用自己定义的。多行 data: 会用换行连接,冒号开头的行是注释,可用于心跳。具体格式见 WHATWG SSE 规范。

浏览器提供了原生客户端:

const source = new EventSource("/jobs/job_001/events");
 
source.addEventListener("content", (event) => {
  const { text } = JSON.parse(event.data);
  console.log(text);
});
 
source.addEventListener("done", () => source.close());
source.addEventListener("failed", () => source.close());

原生 EventSource 负责解析事件和自动重连。需要自定义请求头、明确控制重试或等待业务处理成功后再推进游标时,也可以用 fetch 读取流,但这些能力需要自己实现。下面的完整消费链路使用 fetch。两者接收的是同一种 SSE 格式。MDN 的使用说明给出了原生客户端示例。

此时我们只解决了“怎样持续传消息”。服务端是否保存历史、连接断开后任务是否继续,都还没有答案。

2. 业务建模:任务、事件和连接分别是什么

先定下这个例子的业务要求:一次写作可能持续几分钟,关闭页面不等于取消任务,用户重新打开后能恢复展示。

因此,需要把任务执行和 SSE 连接的生命周期分开:

浏览器 ── POST /jobs ──→ API ──→ jobs + outbox(同一事务)
                                  │
                             Dispatcher
                                  ↓
                                Queue
                                  ↓
                               Worker
                                  │
                       保存 job_events 和结果
                                  ↓
浏览器 ←── SSE Endpoint ←── 按游标读取已提交事件
对象负责什么不由它解决的问题
Job标识一次任务,记录归属、状态和结果位置一次连接是否在线
Worker执行生成,推进任务状态,产生事件浏览器怎样渲染
Event Store持久保存事件,支持按位置重放网络是否成功送达
SSE Endpoint鉴权、读取事件、编码并发送执行长任务
Client应用事件,维护页面状态和消费游标直接操作任务数据库

这里的关键是:Worker 先把事件保存下来,SSE 接口再读取并传给页面。 即使没有浏览器在线,任务也可以继续产生事件。

并不是所有 SSE 接口都需要这些对象。一个只要求当前连接内展示内容、断开就结束的场景,可以更简单。这里引入 Job 和事件日志,是为了满足后台执行与恢复的要求。

先统一前后端的数据契约

type Job = {
  id: string;
  ownerId: string;
  status: "queued" | "running" | "completed" | "failed";
  lastSeq: number;
  resultId?: string;
};
 
type EventBody =
  | { type: "progress"; payload: { percent: number } }
  | { type: "content"; payload: { text: string } }
  | { type: "done"; payload: { resultId: string } }
  | { type: "failed"; payload: { message: string } };
 
type JobEvent = EventBody & {
  jobId: string;
  seq: number;
};
 
type ClientState = {
  jobId: string;
  text: string;
  progress: number;
  lastSeq: number;
  status: "watching" | "completed" | "failed";
  resultId?: string;
  error?: string;
};

本文约定每个任务的 seq 从 1 开始,连续递增,按提交可见顺序出现。SSE 协议中的 ID 本来是字符串,服务端发送时把 seq 转成字符串即可。这是我们的业务约定,并非协议要求。

Job.status 表示后台任务生命周期,ClientState.status 表示页面已消费到的业务结果。网络暂时掉线应另记为连接状态,不能直接把任务改成 failed。

3. 创建任务:先得到 jobId,再建立订阅

前端需要两个接口:

POST /jobs
Content-Type: application/json
Idempotency-Key: <本次提交生成并在重试时复用的键>
 
{"prompt":"写一篇关于 SSE 的文章"}
HTTP/1.1 202 Accepted
Content-Type: application/json
 
{"jobId":"job_001"}

随后请求 GET /jobs/job_001/events?after=0。创建和订阅拆开后,重连只会重新读取已有任务,不会再启动一次生成。

后端创建逻辑的核心是一次事务:

async function createJob(ownerId: string, key: string, prompt: string) {
  return db.transaction(async (tx) => {
    // (ownerId, key) 有唯一约束;同键异参应拒绝。
    // insertOrGet 必须在并发请求下也只创建一个 Job。
    const { job, created } = await tx.jobs.insertOrGet({
      ownerId, key, prompt, status: "queued", lastSeq: 0,
    });
 
    if (created) {
      await tx.outbox.insert({
        key: `enqueue:${job.id}`,
        type: "job.created",
        payload: { jobId: job.id },
      });
    }
    return { jobId: job.id };
  });
}

为什么多一张 Outbox 表?如果先写 Job,再直接向队列发消息,中间进程退出,就会出现“任务创建成功,却没有人执行”。这里把待投递消息和 Job 一起提交,独立 Dispatcher 随后读取 Outbox 并投递队列,失败可以重试。

Dispatcher 也可能在投递成功、标记成功前退出,因此队列仍可能收到重复消息。Worker 必须通过任务领取、租约等机制防止重复执行。Outbox 解决的是投递缺口,不自动保证整个业务只执行一次。

4. Worker:执行任务,保存可重放的事件

Worker 不依赖 SSE 请求。它领取 Job 后调用生成服务,把进度和正文写入事件日志:

async function runJob(jobId: string) {
  const lease = await jobs.claim(jobId);
  if (!lease) return; // 已完成,或正在由其他 Worker 执行
 
  // withLease 负责续租;丢失租约时中止执行。
  await jobs.withLease(lease, async (signal) => {
    await events.append(lease, {
      type: "progress", payload: { percent: 10 },
    });
 
    let text = "";
    for await (const chunk of generateArticle(lease.prompt, signal)) {
      text += chunk;
      await events.append(lease, {
        type: "content", payload: { text: chunk },
      });
    }
 
    await jobs.completeWithResultAndEvent(lease, text);
  });
}

这些存储接口隐藏的是数据库细节,不应隐藏一致性要求:

  • append 在事务中锁住对应 Job 行,校验当前租约,分配下一个 seq,插入事件并更新 lastSeq。
  • completeWithResultAndEvent 在同一事务里保存最终结果、把 Job 改成 completed、追加 done 事件。
  • 不可恢复的执行错误,通过类似事务写入 failed 状态与事件;丢失租约的旧 Worker 无权再写。

单次追加事件可以理解为下面这组 SQL。事务提交后,新游标和新事件一起可见:

BEGIN;
SELECT id FROM jobs WHERE id = $1 FOR UPDATE;
-- 此处校验租约有效,且任务仍允许写入。
UPDATE jobs SET last_seq = last_seq + 1
WHERE id = $1 RETURNING last_seq;
INSERT INTO job_events (job_id, seq, type, payload)
VALUES ($1, $2, $3, $4); -- $2 是上一步返回的 last_seq
COMMIT;

job_events 需要 (job_id, seq) 唯一索引。不能只认为“全局自增 ID 一定按提交顺序可见”:两个并发事务可能先分配小 ID,却后提交,读者推进游标后就会漏掉它。

这段 Worker 代码展示一次正常执行。Worker 中途崩溃后的恢复还需要检查点或新的执行版本;不能把整次生成无条件重跑,再向旧正文继续追加。事件可重放和任务可恢复,是需要分别设计的两件事。

5. SSE 接口:按游标读日志,编码并持续发送

客户端携带 after=3,意思是已经成功应用这个任务的前 3 条事件。服务端查询:

SELECT job_id, seq, type, payload
FROM job_events
WHERE job_id = $1 AND seq > $2
ORDER BY seq ASC
LIMIT 100;

发送前,把业务事件编码成 SSE:

function encodeSse(event: JobEvent): string {
  return `id: ${event.seq}\nevent: ${event.type}\ndata: ${JSON.stringify(event.payload)}\n\n`;
}

下面用一个框架无关的 sink 表示 HTTP 响应。适配层设置 Content-Type: text/event-stream、Cache-Control: no-cache, no-transform,负责及时 flush;sink.write 必须等待背压解除,连接关闭时让 signal 中止。

async function streamJob(
  jobId: string,
  after: number,
  sink: { write(frame: string): Promise<void> },
  signal: AbortSignal,
) {
  let cursor = after;
  let heartbeatAt = Date.now();
 
  while (!signal.aborted) {
    const batch = await events.listAfter(jobId, cursor, 100);
 
    for (const event of batch) {
      await sink.write(encodeSse(event));
      cursor = event.seq;
      if (event.type === "done" || event.type === "failed") return;
    }
 
    if (batch.length > 0) continue;
 
    // 必须先查事件再查终态;终态事务提交后,lastSeq 才确定。
    const job = await jobs.get(jobId);
    if (isTerminal(job.status) && cursor >= job.lastSeq) return;
 
    if (Date.now() - heartbeatAt >= 15_000) {
      await sink.write(": heartbeat\n\n");
      heartbeatAt = Date.now();
    }
    await sleep(500, signal); // 可取消的等待
  }
}

进入此函数前,路由必须完成登录校验、任务归属检查、游标校验,拒绝负数、非安全整数以及超过任务最新位置的游标。HTTP 适配层还要在正常结束时关闭响应,在异常或取消时释放资源。

这里选择周期性读持久日志,避免一开始就引入“历史读完后切换实时订阅”的竞态。浏览器看到的仍是一个 SSE 长连接;周期性查询发生在服务器和数据库之间。

代价是每个活跃连接都有查库压力。规模扩大后,可以增加通知来唤醒读取,并保留周期补查。通知只提示“可能有新事件”,真正的数据仍按游标从日志读取,这样通知丢失也能补回。

还有两个容易混淆的位置:服务端的 cursor 表示这个连接已经写出了哪些事件;客户端的 lastSeq 表示页面已经成功应用到哪里。写进响应不等于客户端已经处理成功。

6. 前端:字节流 → 完整事件 → 业务状态

fetch 的一次 reader.read() 不对应一条 SSE 事件:一个事件可能跨多个 chunk,一次 read 也可能拿到多条事件。因此前端要先解码字节,再缓存文本,按事件边界解析,最后更新业务状态。

下面的解析器只处理本文编码器输出的格式:LF 换行、每条事件一个 JSON data:、显式 id 和 event,以及注释心跳。接第三方 SSE 时,应使用覆盖完整协议的解析器,处理 CRLF、多行 data 等情况。

async function* readEvents(body: ReadableStream<Uint8Array>) {
  const reader = body.getReader();
  const decoder = new TextDecoder();
  let buffer = "";
 
  try {
    while (true) {
      const { value, done } = await reader.read();
      if (done) return; // 未完成的帧不应用,重连时由日志补回
      buffer += decoder.decode(value, { stream: true });
 
      let boundary: number;
      while ((boundary = buffer.indexOf("\n\n")) >= 0) {
        const frame = buffer.slice(0, boundary);
        buffer = buffer.slice(boundary + 2);
        if (frame.startsWith(":")) continue;
 
        const fields = Object.fromEntries(
          frame.split("\n").map((line) => {
            const i = line.indexOf(":");
            return [line.slice(0, i), line.slice(i + 1).replace(/^ /, "")];
          }),
        );
        // 运行时校验 seq 为正安全整数、type 和 payload 匹配;失败就抛错。
        yield validateWireEvent({
          seq: Number(fields.id),
          type: fields.event,
          payload: JSON.parse(fields.data),
        });
      }
    }
  } finally {
    try { await reader.cancel(); }
    finally { reader.releaseLock(); }
  }
}

TextDecoder 的 stream: true 用来保留跨 chunk 的 UTF-8 字节,buffer 用来保留跨 chunk 的事件文本。这是两个不同层面的边界。

解析完的消息还不能直接随意追加。用一个纯函数把内容和消费游标一起推进:

function applyEvent(state: ClientState, event: JobEvent): ClientState {
  if (event.jobId !== state.jobId) throw new Error("任务不匹配");
  if (event.seq <= state.lastSeq) return state;
  if (event.seq !== state.lastSeq + 1) throw new Error("事件存在缺口");
 
  const next = { ...state, lastSeq: event.seq };
  switch (event.type) {
    case "progress":
      return { ...next, progress: event.payload.percent };
    case "content":
      return { ...next, text: state.text + event.payload.text };
    case "done":
      return { ...next, status: "completed", resultId: event.payload.resultId };
    case "failed":
      return { ...next, status: "failed", error: event.payload.message };
  }
}

连接层将这些步骤串起来:

async function consumeOnce(state: ClientState, signal: AbortSignal) {
  const response = await fetch(
    `/jobs/${encodeURIComponent(state.jobId)}/events?after=${state.lastSeq}`,
    { signal, headers: { Accept: "text/event-stream" }, cache: "no-store" },
  );
  if (!response.ok) throw new HttpError(response.status);
  if (!response.headers.get("content-type")?.startsWith("text/event-stream")) {
    throw new Error("响应类型错误");
  }
  if (!response.body) throw new Error("响应体为空");
 
  for await (const wire of readEvents(response.body)) {
    const next = applyEvent(state, { ...wire, jobId: state.jobId });
    // 一个原子检查点:同时保存 text、status、lastSeq 等完整状态。
    // 例如在一个 IndexedDB 事务里写入同一条记录。
    await checkpoints.save(next);
    state = next;
    render(state);
    if (state.status !== "watching") return;
  }
  throw new Error("终态事件到达前,连接结束");
}

外层重连管理器每次从检查点恢复状态,再调用 consumeOnce。网络错误或可重试的服务端错误使用带抖动的退避;鉴权失败、非法数据和检查点保存失败需要停止并处理,不能全部无限重试。校验失败的事件不能推进游标。

页面进入时恢复检查点;没有检查点就从空状态、lastSeq = 0 开始。切换任务或离开页面时,用 AbortController.abort() 关闭订阅,并等待旧消费退出,避免同一页面启动两个消费者。中止订阅只意味着不再观看;取消后台任务需要另一个明确的业务接口。

上面省略了框架组件封装、重试调度器和存储实现,但保留了消费顺序:解析 → 校验 → 应用 → 一致保存状态与游标 → 更新 UI。

7. 断线重连与去重,放回这条链路里理解

假设任务已经产生 1~5,客户端检查点停在 3:

客户端成功保存:正文前半部分 + lastSeq = 3
                       ↓ 断网
Worker 继续保存:4 content、5 done
                       ↓ 重连 after=3
服务端查询:seq > 3,依次发送 4、5
                       ↓
客户端补上正文,保存 completed 状态,停止订阅

能恢复,是因为 Worker 独立执行、服务端保存了事件、客户端保存了对应状态和游标。仅仅给消息加一个 ID,不会自动获得这些能力。

已经只查询后面的事件,为什么还要去重?

正常情况下,客户端携带 3,服务端查询 seq > 3,当然不会再次返回 3。

重复消费的风险来自恢复状态不一致等异常窗口。例如正文中的第 3 段已经保存,但游标保存前页面崩溃,恢复时只能拿到 2,服务端就会正确地再次发送 3。如果直接执行 text += event.payload.text,正文会重复。

因此本文把正文和游标保存为一个原子检查点:恢复出来的要么都是应用之前的状态,要么都是应用之后的状态。在单任务有序串行消费的前提下,seq <= lastSeq 可以跳过已处理事件,不需要再无限增长一个 processedIds 集合。

如果事件会触发付款、发送邮件等外部副作用,这个本地检查点就不够了,必须由实际执行副作用的服务提供幂等机制。

页面刷新后,只恢复游标够不够?

不够。如果恢复游标 3,却没有恢复前 3 条事件构建的正文,从 4 开始就会缺内容。恢复单位必须是“状态 + 对应游标”。也可以从 0 重放完整历史,或先取服务端快照再补增量。

原生 EventSource 自动维护的事件 ID 也不是业务处理成功的确认,不能假定它会等待异步保存检查点结束。本文的 fetch 客户端显式管理这个边界。

8. 上生产前,还要补齐哪些边界

到这里,正常生成和断线恢复已经有了清晰的职责分配。实际部署时,还需要把下面这些行为落实到具体实现和验证中:

边界需要明确的行为
代理与平台禁用流式响应缓冲,验证首条事件能及时抵达;按实际平台配置超时,接受连接可能被定时切断
心跳与背压空闲时发送注释心跳;慢客户端不能让内存无限堆积,限制缓冲并允许断开后重放
事件保留期限旧日志被清理时,不得静默从剩余事件继续;可返回约定的 410,让客户端取得状态与游标一致的快照
终态与结果完成状态、结果和 done 事件一起提交;客户端处理终态后停止重连
Worker 重试租约防止旧执行者继续写;通过检查点或执行版本定义重试,避免正文混入两次生成
权限与隔离每次订阅都校验用户与任务归属;jobId 和游标都不能代替权限
可观测性用 jobId 关联创建、执行、存储、订阅日志,观察事件延迟、重连次数和失败位置

验证时可以沿同一个任务做几次故障注入:接到一半主动断开后重连、刷新后恢复检查点、重复投递一条事件、Worker 完成后再进入页面、用过期游标请求。分别检查正文没有遗漏或重复,任务终态和最终结果一致,客户端不会无休止重连。

回头再看“页面逐段显示文章”,就能沿着同一条链路定位问题:创建接口给出任务身份,Worker 产生业务事实,事件日志保留恢复依据,SSE 接口负责传输,客户端把有序事件变成页面状态。以后遇到丢内容、重复内容或断线后停住,都可以先问:问题发生在产生、保存、发送、消费,还是恢复位置的那一步?

目录 · 收起