流式响应的断点续传:前端怎么记住读到哪一行了

断点续传的核心不是“记住行号”,而是用服务端返回的事件 ID 作为断点位置——前端每次消费一行,就把最后成功处理的 event ID 存下来,重连时通过 HTTP 的 Last-Event-ID 请求头或自定义参数告诉服务端从哪里继续推。行号是客户端自己的渲染索引,不可靠;事件 ID 才是断点续传的正确锚点。

下面我会拆解一整套生产级方案,包括 SSE 协议里 id 字段的正确用法、前端消费队列与持久化策略、重连时的状态恢复,以及服务端配合的最小改造点。

SSE 的 id 字段就是天然断点标记

很多人用 SSE 只关注 data 字段,完全忽略了 id。规范里写得很清楚:服务端可以在每个事件前加一行 id: xxx,浏览器原生 EventSource 会自动把这个值存下来,下次重连时自动带上 Last-Event-ID 请求头。

一个符合规范的 SSE 流长这样:

id: evt_001
data: {"token": "Hello", "index": 0}

id: evt_002
data: {"token": " world", "index": 1}

id: evt_003
data: {"token": "!", "index": 2}

前端如果用原生 EventSource,重连是自动的,Last-Event-ID 也是自动带上的。但大模型 API 的流式响应通常不是标准 SSE 的 text/event-stream,而是 text/plain 里嵌着一行行 data: {...} 的 JSON 流(OpenAI、Anthropic 都是这样),而且大多是 POST 请求,EventSource 只支持 GET。所以实际场景里我们基本都用 fetch + ReadableStream 手动解析。

手动解析就意味着,id 的提取和断点续传逻辑都得自己写。

手动解析流时,如何提取并记录事件 ID

先看一个基础的手动流解析实现:

async function streamChat(
  body: object,
  lastEventId: string | null,
  onToken: (token: string, eventId: string) => void
) {
  const headers: Record<string, string> = {
    'Content-Type': 'application/json',
  };
  if (lastEventId) {
    headers['X-Last-Event-ID'] = lastEventId;
  }

  const response = await fetch('/api/chat/stream', {
    method: 'POST',
    headers,
    body: JSON.stringify(body),
  });

  const reader = response.body!.getReader();
  const decoder = new TextDecoder();
  let buffer = '';
  let currentEventId = '';

  while (true) {
    const { done, value } = await reader.read();
    if (done) break;

    buffer += decoder.decode(value, { stream: true });
    const lines = buffer.split('\n');
    // 保留最后一个可能不完整的行
    buffer = lines.pop() || '';

    for (const line of lines) {
      if (line.startsWith('id: ')) {
        currentEventId = line.slice(4).trim();
      } else if (line.startsWith('data: ')) {
        const data = line.slice(6);
        if (data === '[DONE]') return;

        try {
          const parsed = JSON.parse(data);
          const token = parsed.choices?.[0]?.delta?.content ?? '';
          if (token) {
            onToken(token, currentEventId);
          }
        } catch {
          // 某些行可能不是完整 JSON,忽略
        }
      }
    }
  }
}

这里的关键操作是:每解析出一个 id: 行就更新 currentEventId,然后在触发 onToken 回调时把当前的事件 ID 一并传出去。这样消费方就能在拿到 token 的同时知道它对应的事件 ID,后续持久化时就能精确标记“这个 token 我已经消费过了”。

消费端:用一个确认队列保证不丢不重

断点续传最怕两件事:丢 token(少字)和重 token(重复字)。如果不做任何去重,重连后服务端从某个事件 ID 开始重新推送,前端可能会把已经渲染过的 token 再渲染一遍。

解决方案是维护一个 已确认消费的事件 ID 集合,渲染前做去重:

class StreamConsumer {
  private consumedIds = new Set<string>();
  private lastPersistedId: string | null = null;
  private renderedText = '';
  private pendingPersist: ReturnType<typeof setTimeout> | null = null;

  constructor(
    private persistKey: string,
    private onUpdate: (text: string) => void
  ) {
    // 从 localStorage 恢复上次消费到哪个事件 ID
    const saved = localStorage.getItem(this.persistKey);
    if (saved) {
      try {
        const state = JSON.parse(saved);
        this.lastPersistedId = state.lastEventId;
        this.renderedText = state.text;
        // 恢复时不需要重建 consumedIds,因为 lastPersistedId 之前的都确认消费了
        this.consumedIds.add(state.lastEventId);
        this.onUpdate(this.renderedText);
      } catch {}
    }
  }

  consume(token: string, eventId: string) {
    // 去重:如果这个事件 ID 已经消费过,跳过
    if (this.consumedIds.has(eventId)) return;

    this.consumedIds.add(eventId);
    this.renderedText += token;
    this.onUpdate(this.renderedText);

    // 防抖持久化,避免频繁写 localStorage
    if (this.pendingPersist) clearTimeout(this.pendingPersist);
    this.pendingPersist = setTimeout(() => {
      this.lastPersistedId = eventId;
      localStorage.setItem(
        this.persistKey,
        JSON.stringify({
          lastEventId: eventId,
          text: this.renderedText,
          timestamp: Date.now(),
        })
      );
    }, 200);
  }

  getLastEventId(): string | null {
    return this.lastPersistedId;
  }

  flush() {
    if (this.pendingPersist) {
      clearTimeout(this.pendingPersist);
      this.pendingPersist = null;
    }
    // 最终持久化
    const lastId = [...this.consumedIds].pop() || null;
    if (lastId) {
      localStorage.setItem(
        this.persistKey,
        JSON.stringify({
          lastEventId: lastId,
          text: this.renderedText,
          timestamp: Date.now(),
        })
      );
    }
  }
}

consume 方法做三件事:用 Set 去重、拼接文本、防抖写 localStorage。去重是防御性的——正常流程下服务端从正确的事件 ID 之后开始推送,不应该有重复,但网络异常、服务端实现有 bug 时,这个 Set 就是最后一道防线。

防抖周期我设 200ms,这个值权衡了持久化频率和崩溃丢数据的风险。200ms 意味着最坏情况下丢失最后 200ms 内产生的 token,在 60 tokens/s 的生成速度下大约是 12 个 token,对用户体验影响很小。如果你的场景对数据完整性要求极高,可以降到 50ms 甚至同步写,但要评估 localStorage 的写入开销。

重连时如何把断点传给服务端

前端存好了 lastEventId,重连时要把它传给服务端。我习惯用自定义请求头 X-Last-Event-ID,因为这样不污染请求体,服务端也容易在中间件层统一处理:

async function startOrResumeChat(sessionId: string, prompt: string) {
  const consumer = new StreamConsumer(
    `chat:${sessionId}`,
    (text) => {
      // 更新 UI
      document.getElementById('output')!.textContent = text;
    }
  );

  const lastEventId = consumer.getLastEventId();

  await streamChat(
    { sessionId, prompt },
    lastEventId,
    (token, eventId) => {
      consumer.consume(token, eventId);
    }
  );

  // 流正常结束,最终持久化
  consumer.flush();
}

服务端收到 X-Last-Event-ID 后,需要从该事件 ID 的下一个事件开始推送。如果服务端是自建的,最简单的做法是把整个对话的所有事件持久化存储(比如按 sessionId 存一个事件列表),然后根据 lastEventId 找到起始位置:

# 伪代码,服务端逻辑
last_event_id = request.headers.get('X-Last-Event-ID')
events = get_events_for_session(session_id)

if last_event_id:
    # 找到 last_event_id 的索引,从下一个开始
    start_idx = next(
        (i for i, e in enumerate(events) if e.id == last_event_id),
        -1
    ) + 1
else:
    start_idx = 0

# 把 start_idx 之后的事件逐个推送给客户端

如果用的是 OpenAI 或 Anthropic 的 API,它们本身不支持从指定事件 ID 续传。这种情况下你需要在服务端做一个代理层,把 API 返回的流式数据先完整缓存下来,同时给每个 token 分配递增的事件 ID,然后按上面的逻辑支持续传。缓存策略可以用 Redis,key 是 session_id,value 是一个有序列表,每个元素包含 eventIdtokentimestamp

处理流中断的几种场景

网络断开、用户切换 App、浏览器崩溃——这三种情况需要区别对待。

网络断开fetchReadableStream 会抛出 TypeError: network error。在 streamChat 外层包一个重试循环,指数退避,最大重试 5 次:

async function streamWithRetry(
  body: object,
  consumer: StreamConsumer,
  maxRetries = 5
) {
  for (let attempt = 0; attempt < maxRetries; attempt++) {
    try {
      await streamChat(
        body,
        consumer.getLastEventId(),
        (token, eventId) => consumer.consume(token, eventId)
      );
      return; // 正常结束
    } catch (err) {
      if (attempt === maxRetries - 1) throw err;
      const delay = Math.min(1000 * Math.pow(2, attempt), 16000);
      await new Promise((resolve) => setTimeout(resolve, delay));
    }
  }
}

用户切换 App(页面隐藏):移动端 Safari 和部分 Android WebView 会在页面进入后台时暂停 requestAnimationFrame 和某些定时器,但 fetch 流通常不会中断。不过如果后台时间过长(超过 30 秒),系统可能会杀掉网络连接。用 document.visibilitychange 监听页面可见性变化,在页面恢复可见时检查流是否还活着:

document.addEventListener('visibilitychange', () => {
  if (document.visibilityState === 'visible') {
    // 检查 reader 是否已被取消或出错
    // 如果流已断,触发重连
  }
});

浏览器崩溃:这是最简单的情况——崩溃时 localStorage 里已经持久化了最后的事件 ID 和已渲染文本。用户重新打开页面时,StreamConsumer 构造函数会自动从 localStorage 恢复状态,UI 立即显示之前的文本,同时向服务端发起续传请求,从断点继续生成后续内容。整个过程对用户来说就像什么都没发生过。

常见问题

为什么不用行号或 token 索引做断点标记?

行号只在客户端当前连接内有效,重连后服务端可能因为网络缓冲、重试等原因推送的行边界与上次不一致。token 索引同理,它是你本地拼接的顺序号,服务端根本不认识。事件 ID 是服务端和客户端共同认可的、全局唯一的标记,双方对“这个事件之后”有完全一致的理解。

如果服务端不支持 SSE 的 id 字段怎么办?

在代理层手动注入。你可以在服务端收到大模型 API 的流式响应后,给每个 data: 行前面插入一个 id: 行,用自增序号或 UUID 都可以。例如第一个 token 对应 id: 1,第二个对应 id: 2,简单直接。关键是保证在一个 session 内事件 ID 全局唯一且有序。

localStorage 的 5MB 限制够用吗?

一个对话的断点状态通常只有几百字节(一段文本加一个事件 ID),完全够用。但如果你要支持多轮对话,建议按 sessionId 分 key 存储,每次只恢复当前对话的状态,不要把所有历史对话都塞进一个 key 里。如果文本特别长(几万字),可以考虑用 IndexedDB 替代 localStorage,但绝大多数场景没必要。

断点续传后,前端怎么保证 UI 不闪烁?

恢复时先立即渲染 localStorage 中保存的完整文本,这个操作是同步的,用户打开页面瞬间就能看到之前的内容。然后异步发起续传请求,新 token 到达后直接追加到已有文本后面。整个过程 UI 只做增量更新,不会出现“清空再重填”的闪烁。