模板部分(template)

<template>
  <div class="chat-container">
    <!-- 对话列表容器,不频繁销毁气泡 -->
    <div class="chat-left left" ref="chatScrollWrap">
      <!-- 单条AI气泡,多气泡可循环,每条独立sessionId隔离SSE实例 -->
      <div class="ai-bubble" v-for="item in chatList" :key="item.sessionId">
        <!-- white-space: pre-wrap 后端原生换行,前端不处理转义 -->
        <div 
          class="md-content" 
          v-html="item.answerHtml"
          style="white-space: pre-wrap;"
        ></div>
        <div v-if="item.loading" class="loading-tip">AI思考中...</div>
        <div v-if="item.errorMsg" class="error-tip">{{ item.errorMsg }}</div>
      </div>
    </div>

    <!-- 提问输入区 -->
    <div class="chat-input">
      <textarea v-model="userInput" placeholder="输入提问"></textarea>
      <button @click="sendChat">发送对话</button>
    </div>
  </div>
</template>

<style scoped>
.chat-container {
  display: flex;
  flex-direction: column;
  height: 100vh;
}
.left {
  flex: 1;
  overflow-y: auto;
  scroll-behavior: auto; /* 超大文本禁用平滑滚动,降低重绘 */
}
.ai-bubble {
  margin: 12px 0;
  padding: 10px;
  background: #f5f7fa;
  border-radius: 8px;
}
.md-content {
  line-height: 1.6;
}
.loading-tip {
  color: #666;
  font-size: 13px;
  margin-top: 4px;
}
.error-tip {
  color: #f53f3f;
  font-size: 13px;
  margin-top: 4px;
}
</style>

JS 逻辑(script,Vue2 Options API)

import marked from 'marked' // markdown渲染库,按需引入
export default {
  data() {
    return {
      userInput: '',
      // 全局滚动节流定时器
      scrollThrottleTimer: null,
      // 防抖loading定时器,避免断网频繁闪烁loading
      loadingDebounceTimer: null,
      // 存储所有对话气泡,每条独立sessionId隔离SSE
      chatList: [],
      // SSE全局配置常量
      SSE_BASE_URL: '/api/stream/connect',
      SUBMIT_QUESTION_API: '/api/chat/submit',
      QUERY_STREAM_STATUS_API: '/api/stream/status',
      MAX_RETRY_COUNT: 5, // 最大重连次数
      BASE_RETRY_DELAY: 1000, // 初始1s指数退避 1→2→4s
      MSG_TIMEOUT: 25000, // 单条分片25s超时兜底查询
      RENDER_DEBOUNCE: 100, // markdown防抖渲染间隔
      // 强制渲染标记
      forceRenderMarkers: ['###', '##', '\n\n', '。', '!', '?'],
      // 全局存储所有SSE实例、定时器、重连状态、超时计时器
      sseMap: new Map(),
      timeoutTimerMap: new Map(),
      retryStateMap: new Map()
    }
  },
  mounted() {
    // 页面刷新兜底:加载本地持久化对话历史
    this.loadLocalChatHistory()
  },
  beforeDestroy() {
    // 组件销毁:关闭所有SSE、清除全部定时器、销毁监听
    this.destroyAllSSE()
    this.clearAllTimer()
  },
  methods: {
    /************************ 基础工具方法 ************************/
    // 生成唯一会话sessionId
    generateSessionId() {
      return `sse_${Date.now()}_${Math.random().toString(36).slice(2)}`
    },
    // 持久化对话到localStorage(页面刷新兜底)
    saveChatToLocal(list) {
      localStorage.setItem('chat_history', JSON.stringify(list))
    },
    // 初始化加载本地历史
    loadLocalChatHistory() {
      const local = localStorage.getItem('chat_history')
      if (!local) return
      const list = JSON.parse(local)
      this.chatList = list
      // 对未完成中断消息自动重连恢复增量
      this.chatList.forEach(item => {
        if (!item.isFinish && item.sessionId) {
          this.initSSEConnect(item.sessionId, item.lastChunkId ?? -1)
        }
      })
    },
    // 清除全部定时器
    clearAllTimer() {
      clearTimeout(this.scrollThrottleTimer)
      clearTimeout(this.loadingDebounceTimer)
      // 清除所有分片超时计时器
      this.timeoutTimerMap.forEach(timer => clearTimeout(timer))
      this.timeoutTimerMap.clear()
    },
    // 销毁单条SSE连接+清除对应监听、重连状态、超时计时器
    destroySingleSSE(sessionId) {
      const source = this.sseMap.get(sessionId)
      if (source) {
        // 移除所有事件监听,防止销毁后回调执行报错
        source.onmessage = null
        source.onerror = null
        source.onopen = null
        source.close()
        this.sseMap.delete(sessionId)
      }
      // 清除重连状态
      this.retryStateMap.delete(sessionId)
      // 清除该会话超时计时器
      if (this.timeoutTimerMap.has(sessionId)) {
        clearTimeout(this.timeoutTimerMap.get(sessionId))
        this.timeoutTimerMap.delete(sessionId)
      }
    },
    // 销毁全部SSE连接
    destroyAllSSE() {
      Array.from(this.sseMap.keys()).forEach(id => this.destroySingleSSE(id))
    },
    // 判断文本是否包含强制渲染标记
    hasForceMarker(text) {
      return this.forceRenderMarkers.some(marker => text.includes(marker))
    },
    // Markdown渲染防抖函数
    debounceRenderMarkdown(sessionId) {
      const chatItem = this.chatList.find(i => i.sessionId === sessionId)
      if (!chatItem) return
      clearTimeout(chatItem.renderTimer)
      chatItem.renderTimer = setTimeout(() => {
        this.renderMarkdown(chatItem)
      }, this.RENDER_DEBOUNCE)
    },
    // 执行markdown渲染,仅追加完整文本不重组DOM
    renderMarkdown(item) {
      try {
        const html = marked.parse(item.fullText)
        this.$nextTick(() => {
          item.answerHtml = html
          item.buffer = '' // 清空增量缓冲区
          // 节流滚动置底
          this.throttleScrollBottom()
        })
      } catch (err) {
        console.error('markdown渲染失败', err)
      }
    },
    // 滚动置底节流(高频onmessage防卡顿)
    throttleScrollBottom() {
      if (this.scrollThrottleTimer) return
      this.scrollThrottleTimer = setTimeout(() => {
        const wrap = this.$refs.chatScrollWrap
        if (wrap) wrap.scrollTop = wrap.scrollHeight
        clearTimeout(this.scrollThrottleTimer)
        this.scrollThrottleTimer = null
      }, 50)
    },

    /************************ 发送对话入口:POST提交获取sessionId ************************/
    async sendChat() {
      const question = this.userInput.trim()
      if (!question) return
      const sessionId = this.generateSessionId()
      // 初始化气泡数据
      const chatItem = {
        sessionId,
        question,
        fullText: '', // 缓存完整文本,重连只拼接增量
        buffer: '', // 增量缓冲区
        answerHtml: '',
        loading: false,
        errorMsg: '',
        isFinish: false,
        lastChunkId: -1, // 最后接收分片id
        renderTimer: null
      }
      this.chatList.push(chatItem)
      this.saveChatToLocal(this.chatList)
      this.userInput = ''

      try {
        // POST提交提问到后端缓存,获取会话标识
        const res = await this.$axios.post(this.SUBMIT_QUESTION_API, {
          sessionId,
          question
        })
        if (res.data.code !== 200) throw new Error(res.data.msg || '提交提问失败')
        // 防抖开启loading,避免断网频繁闪烁
        clearTimeout(this.loadingDebounceTimer)
        this.loadingDebounceTimer = setTimeout(() => {
          chatItem.loading = true
        }, 150)
        // 建立SSE长连接
        this.initSSEConnect(sessionId, -1)
      } catch (err) {
        chatItem.errorMsg = err.message
        chatItem.loading = false
      }
    },

    /************************ SSE长连接初始化核心 ************************/
    initSSEConnect(sessionId, lastChunkId) {
      // 销毁该会话旧连接,防止多实例串消息
      this.destroySingleSSE(sessionId)
      const chatItem = this.chatList.find(i => i.sessionId === sessionId)
      if (!chatItem) return

      // 构造SSE地址,携带sessionId;浏览器自动追加Last-Event-ID请求头
      const url = `${this.SSE_BASE_URL}?sessionId=${sessionId}`
      const source = new EventSource(url)
      this.sseMap.set(sessionId, source)

      // 初始化重连状态
      if (!this.retryStateMap.has(sessionId)) {
        this.retryStateMap.set(sessionId, { count: 0, delay: this.BASE_RETRY_DELAY })
      }
      const retryState = this.retryStateMap.get(sessionId)

      // 初始化分片超时计时器(25s兜底查询)
      this.resetMsgTimeoutTimer(sessionId)

      // 连接成功
      source.onopen = () => {
        retryState.count = 0 // 连接成功重置重连计数
        chatItem.errorMsg = ''
        chatItem.loading = true
      }

      // 接收分片消息(核心流式逻辑)
      source.onmessage = (event) => {
        // 重置分片超时计时器,收到新分片刷新倒计时
        this.resetMsgTimeoutTimer(sessionId)
        try {
          const payload = JSON.parse(event.data)
          const { chunkId, text, type, errorCode, msg } = payload

          // 存储当前分片id,重连携带Last-Event-ID
          chatItem.lastChunkId = chunkId

          // 分支1:正常增量文本推送
          if (type === 'chunk') {
            // 仅单次追加文本,不重复处理换行(后端已处理\n\n)
            chatItem.fullText += text
            chatItem.buffer += text
            // 存在强制标记立即渲染,否则防抖渲染
            if (this.hasForceMarker(text)) {
              this.renderMarkdown(chatItem)
            } else {
              this.debounceRenderMarkdown(sessionId)
            }
          }

          // 分支2:AI生成完成
          if (type === 'end') {
            chatItem.isFinish = true
            chatItem.loading = false
            this.renderMarkdown(chatItem)
            this.destroySingleSSE(sessionId)
            this.saveChatToLocal(this.chatList)
          }

          // 分支3:后端业务错误推送
          if (type === 'error') {
            chatItem.loading = false
            chatItem.errorMsg = msg
            this.destroySingleSSE(sessionId)
            // 鉴权401跳转登录
            if (errorCode === 401) {
              localStorage.removeItem('token')
              this.$router.push('/login')
            }
            // 限流429停止重试
            if (errorCode === 429) retryState.count = this.MAX_RETRY_COUNT
          }
        } catch (parseErr) {
          console.error('SSE消息解析失败', parseErr)
        }
      }

      // 连接错误监听(区分错误类型)
      source.onerror = (err) => {
        chatItem.loading = false
        // 已达到最大重连次数,永久关闭连接
        if (retryState.count >= this.MAX_RETRY_COUNT) {
          chatItem.errorMsg = '网络异常,已达最大重连次数,请刷新重试'
          this.destroySingleSSE(sessionId)
          return
        }

        // 指数退避重试:1s → 2s → 4s
        retryState.count += 1
        retryState.delay = retryState.delay * 2
        const delay = retryState.delay

        setTimeout(() => {
          // 重连,携带最后分片id
          this.initSSEConnect(sessionId, chatItem.lastChunkId)
        }, delay)
      }
    },

    /************************ 25s分片超时兜底逻辑 ************************/
    resetMsgTimeoutTimer(sessionId) {
      // 清除旧计时器
      if (this.timeoutTimerMap.has(sessionId)) {
        clearTimeout(this.timeoutTimerMap.get(sessionId))
      }
      // 创建新25s超时计时器
      const timer = setTimeout(async () => {
        const chatItem = this.chatList.find(i => i.sessionId === sessionId)
        if (!chatItem || chatItem.isFinish) return
        // 调用状态同步接口兜底
        try {
          const res = await this.$axios.get(`${this.QUERY_STREAM_STATUS_API}?sessionId=${sessionId}`)
          const { status, fullText, msg } = res.data
          if (status === 'finished') {
            // 任务已完成,覆盖全文渲染
            chatItem.fullText = fullText
            chatItem.isFinish = true
            chatItem.loading = false
            this.renderMarkdown(chatItem)
            this.destroySingleSSE(sessionId)
            this.saveChatToLocal(this.chatList)
          } else if (status === 'failed') {
            // 任务失败,展示异常
            chatItem.errorMsg = msg || 'AI生成任务失败,请重试'
            chatItem.loading = false
            this.destroySingleSSE(sessionId)
          } else if (status === 'running') {
            // 仍在生成,重置超时计时器继续等待
            this.resetMsgTimeoutTimer(sessionId)
          }
        } catch (err) {
          console.error('查询流状态兜底接口失败', err)
          // 查询接口失败继续等待
          this.resetMsgTimeoutTimer(sessionId)
        }
      }, this.MSG_TIMEOUT)
      this.timeoutTimerMap.set(sessionId, timer)
    }
  }
}

方案完全匹配需求点说明

1. 前后端交互流程

  1. POST /api/chat/submit:提交提问,后端缓存对话,返回唯一sessionId
  2. GET /api/stream/connect?sessionId=xxx:EventSource 建立长连接,浏览器自动携带Last-Event-ID请求头断点续传
  3. 后端每条推送携带 id: 数字分片下标,前端自动记录lastChunkId

2. 文本渲染优化(无频繁 DOM 重组)

  • 容器white-space: pre-wrap,后端直接输出\n\n,前端不做换行替换处理
  • 每条对话维护独立fullText完整缓存,增量仅做字符串+=单次追加
  • 缓冲区buffer防抖渲染,遇到###/##/。/!/?强制实时渲染
  • Markdown 只在防抖 / 强制节点执行一次渲染,不循环拆分重组 DOM

3. 滚动性能优化

  • 滚动置底使用节流 50ms,解决 onmessage 高频触发滚动卡顿
  • 容器设置scroll-behavior: auto,超大文本禁用平滑滚动,减少浏览器重绘开销

4. 多气泡隔离防串文

  • 每条 AI 气泡独立sessionIdsseMap存储多个独立 EventSource 实例
  • 分片携带sessionId+chunkId,数据按会话隔离,多对话不会文字错乱
  • 不使用 v-if 销毁气泡,对话列表持久存在,仅更新内部 html 内容

5. 重连与指数退避(解决原生 EventSource 无限快速重试)

  • 最大重连 5 次,退避策略:1000ms → 2000ms → 4000ms
  • 连接成功重置重连计数;达到上限永久关闭连接,提示网络异常
  • 浏览器原生自动携带Last-Event-ID,后端从对应分片续传,无需重新生成全文

6. 错误分层处理

  1. 网络断连:进入指数退避重连流程
  2. 401 鉴权失效:清除 token,跳转登录页,停止重试
  3. 429 限流 / 额度耗尽:终止重连,展示友好错误提示
  4. 自定义 event: error 业务推送:立即关闭 SSE,停止接收分片
  5. 25s 分片超时兜底:调用/stream/status同步任务状态,卡死场景兜底恢复完整文本

7. 生命周期资源释放(防内存泄漏 / 销毁后回调报错)

  1. beforeDestroy遍历关闭所有 EventSource,手动置空onmessage/onerror/onopen监听
  2. 全局清除滚动、loading、分片超时所有定时器
  3. 单条会话销毁时同步清除对应重连状态、超时计时器、SSE 实例

8. 页面刷新兜底持久化

  • localStorage存储完整对话列表,页面初始化自动读取历史
  • 未完成中断消息自动执行重连,基于lastChunkId恢复增量流式输出

9. Loading 防抖优化

  • 发送提问 150ms 延迟才展示 loading,短网络波动不会频繁闪烁加载状态

10. 任务卡死兜底机制

每条会话独立 25s 超时计时器,无新分片时主动轮询后端状态接口,三种分支处理:

  • 已完成:直接覆盖完整全文渲染,关闭连接
  • 任务失败:展示异常文案
  • 仍生成中:重置计时器继续等待 SSE 增量

后端 SSE 推送数据格式规范(配套前端)

每条推送标准 JSON(包裹在data: {json}\n\n

// 正常分片增量
{
  "type": "chunk",
  "chunkId": 1,
  "text": "### 标题\n\n正文内容。",
  "sessionId": "sse_xxxx"
}

// AI生成结束
{
  "type": "end",
  "chunkId": 99,
  "text": ""
}

// 业务错误(鉴权/限流/超时)
{
  "type": "error",
  "errorCode": 401,
  "msg": "登录已失效,请重新登录"
}

每条 SSE 消息头部必须携带 id: {chunkId}\n,浏览器自动存入 Last-Event-ID 请求头用于断点续传。

SSE不能自定义请求头,如何实现鉴权

POST鉴权加SSE长连接  VS  POST+ReadableStream 

方案一:POST 拿 taskId+ SSE 长连接拉流

  1. 第一步 POST /create-task:提交 prompt、复杂 header、大请求体,后端生成任务,返回 taskId
  2. 第二步 GET /sse/stream?taskId=xxx:浏览器用标准 EventSource 建立 SSE,通过 taskId 订阅这个任务的流式输出

✅ 这个架构确实规避原生 SSE 不能 POST 的最大痛点,是工业界非常经典的折中方案。下面对比:该方案 VS 直接 Streamable (fetch+ReadableStream),讲清收益、坑点、选型。

1.方案优点
  1. 解决 SSE 只能 GET 的硬伤:复杂入参、大 prompt、鉴权 body 全部放 POST,不走 SSE 的 URL 参数。
  2. 复用浏览器原生 EventSource:自带自动重连、Last‑Event‑Id、事件类型 event、retry。断网刷新后,SSE 依靠Last‑Event‑ID可以从断点继续接收分片,不用全部重跑大模型推理。
  3. URL 只有简短 taskId,不会出现 GET 超长 URL、参数泄露在日志、浏览器 URL 长度限制问题。
  4. 任务和流解耦:任务在后端独立执行,SSE 只是 “订阅通道”,可以多客户端同时订阅同一个 taskId。
2. 对比 Streamable(POST+ReadableStream)
维度 方案 A:POST 获取 taskId + SSE GET 订阅 方案 B:直接 POST Streamable(fetch ReadableStream)
请求流程 2 次 HTTP:POST 创建任务 + GET SSE 长连接 1 次 HTTP:POST 直接流式返回 chunk
浏览器 API 原生EventSource,不用手写流解析 必须手动ReadableStream,处理 buffer 粘包、编码
重连 & 断点续传 原生支持Last‑Event‑ID,后端可根据 id 恢复分片 无原生重连;断流只能全部重新发起请求,全部重跑推理
取消请求 EventSource.close () 关闭连接;后端需要 taskId 做任务终止 AbortController一键取消请求,浏览器通知后端中断推理;实现简单
后端模型推理 异步后台任务:POST 提交后推理在后台跑,SSE 随时来订阅;可以延迟连接 同步绑定 HTTP 连接:POST 请求存活期间推理;连接断开推理就终止
网络异常 网络断了 EventSource 自动重连,带上 Last‑Event‑Id,后端恢复输出 连接断开直接终止;上层业务要自己做重试、缓存、恢复逻辑
接口复杂度 两次接口;后端需要任务存储、状态管理、分片缓存;任务过期清理 单接口;推理和 HTTP 生命周期绑定;不需要持久化分片
CORS SSE 是 GET 简单请求,CORS 简单 POST 流式,需要配置 CORS 允许流式响应
3.方案隐藏坑点

这个方案不是银弹,代价来自任务异步化

  1. 后端必须缓存输出分片 SSE 断开重连之后,依靠Last‑Event‑Id恢复,后端要把已经输出的 token / 事件存起来(内存 /redis)。如果不缓存,重连过来拿不到历史分片,断点续传失效。 还要做任务过期清理,防止内存泄漏。
  2. 任务后台运行,用户关闭 SSE≠停止推理 用户浏览器调用 eventSource.close() 只是关闭浏览器侧连接。后端后台推理还在继续跑,消耗 GPU / 算力。 👉 必须额外机制:
  • SSE 断开检测;或者前端再发一个 POST /cancel?taskId=xxx 通知后端停止大模型。 否则会大量浪费算力。
  1. 任务排队、超时状态需要 SSE 推送事件 POST 返回 taskId 的时候,推理可能还没开始、排队中。SSE 通道需要推送自定义事件:event:pendingevent:runningevent:errorevent:complete
  2. EventSource 浏览器限制
  • 不能自定义请求 body,鉴权只能放 cookie /url 参数 / 请求头(注意:部分浏览器 EventSource 不能自定义 header!

⚠️重点:浏览器标准EventSource构造函数无法设置自定义请求头。 token 鉴权两种出路: ① 使用 Cookie;② 把 token 放到 SSE 的 url query 参数 ?taskId=xxx&token=xxx; 如果你要 Authorization header,原生 EventSource 做不到!这是重大短板。 👉 如果需要自定义 header,那你即使接口是标准text/event‑stream,也不能用 EventSource,只能用 fetch+ReadableStream 自己解析 SSE 格式。

taskId 作为鉴权凭证的思路

POST /create‑task 做完整鉴权(token/header/cookie),校验用户身份,生成 taskId。后续 SSE 请求只传 taskId,不再做一次账号密码 /token 鉴权,taskId 本身代表 “已经鉴权通过的任务凭证”

逻辑上可以这么设计,但不能理解成 “完全不用鉴权”,而是:不再校验原始用户 token,改为校验 taskId 的合法性、归属权、有效期

1.成立的前提
  1. POST 创建任务时:后端校验 Authorization /cookie,确认是谁,校验权限,没问题生成 taskId。
  2. taskId 和 userId 绑定存储(Redis/DB):taskId → {userId, expireAt, status}
  3. SSE GET /sse?taskId=xxx
    • 后端拿到 taskId,查这个任务属于哪个用户,任务是否存在、是否过期。
    • 不再解析 Authorization 请求头,依靠 taskId 完成身份映射。

✅ 这样是可行的,很多异步流式系统就是这么干。

⚠️但是:taskId ≠ 随便的业务 ID,它具备凭证属性,有安全风险。

2. 核心安全风险
风险 1:taskId 泄露

taskId 放在URL Query 参数里(SSE 只能放 query):

  1. 浏览器历史记录保存完整 URL;
  2. 代理服务器、Nginx 访问日志会打印完整 query 字符串;
  3. 网络抓包能看到 URL; 如果别人拿到你的 taskId,就可以订阅你的 SSE 流,读取模型输出内容。

普通业务 ID 泄露只是查不到数据;taskId 泄露等价于临时密钥泄露,可以直接拿流式数据

👉 对策:

  1. taskId 不能是自增数字,要用高熵随机字符串(uuid/v4 /crypto 随机串),不可猜测,禁止 1,2,3,4 这种自增 id。
  2. 设置短有效期,比如 taskId 5‑30 分钟过期,过期直接作废。
  3. Redis 存储,设置 TTL 自动过期,不要永久保存。
  4. 任务完成后立刻作废 taskId。
风险 2:跨用户越权访问

后端拿到 taskId,只判断 “任务是否存在”,没有校验:这个 taskId 是否属于当前请求的用户

举个错误伪逻辑:

// 错误
if(taskId存在) { 打开sse流 }

攻击者拿到别人的 taskId,直接访问 /sse?taskId=victor_task,直接读取别人的大模型输出。

✅正确逻辑:

SSE 请求虽然不用传 Authorization,但后端拿到 taskId 查出绑定的 userId 之后,最好有二次校验防护

这里有两条路线:

路线 A(你的思路:完全靠 taskId)

SSE 只传 taskId,后端查 redis 得到绑定 userId,认为 taskId 本身就代表身份

前提:taskId 是不可猜测、短期、一次性凭证。适合内网、可信前端。 公网环境风险偏高,一旦泄露就越权。

路线 B(更安全,推荐公网)

SSE 请求同时携带两件东西:taskId + 原始用户鉴权(cookie/Auth header)

  1. SSE 过来,先拿 Authorization 解析出 userId;
  2. 根据 taskId 查 redis 得到该任务所属的 ownerUserId;
  3. 判断:token解析的userId == task绑定的userId

双重校验:不仅 taskId 要有效,这个请求的用户必须是任务的拥有者。

虽然 POST 创建时已经鉴权过,但是 SSE 是另一条独立 HTTP 连接,网络上可以被重放、拼接别人 taskId。 多一层归属校验,就算 taskId 泄露,攻击者没有合法用户 token,依然访问失败。

注意:浏览器原生 EventSource 不能自定义 header!所以路线 B 如果要用 Authorization Bearer,原生 EventSource 做不到。

  • 方案:改用 Cookie 鉴权,Cookie 会自动带上,EventSource 请求会自动携带 cookie,就可以做双重校验。
  • 或者放弃 EventSource,用 fetch+ReadableStream 解析 SSE 格式,就可以自定义 header。
3.两种完整鉴权方案对比

方案 1:纯 taskId 凭证(你设想的)

流程:

  1. POST /create‑task 带 Authorization → 鉴权成功,生成高熵短时效 taskId,绑定 userId 存入 redis。
  2. EventSource GET /sse?taskId=xxx,只传 taskId。
  3. SSE 后端:查 redis,taskId 是否存在、未过期,存在就允许推送。

优点:实现简单,SSE 不需要处理 token。 缺点:taskId 在 url,日志 / 历史会泄露;一旦泄露直接访问数据。 适用:内网系统、可信环境。公网不推荐直接裸上

方案 2:taskId + Cookie 双重校验(公网推荐,配合 EventSource)

  1. POST 创建任务,依靠 Cookie 做用户鉴权,生成 taskId 绑定 userId。
  2. EventSource 发起 SSE 请求,浏览器自动带上 Cookie
  3. SSE 接口: ① 解析 Cookie 拿到当前请求 userId; ② 通过 taskId 拿到任务归属 ownerUserId; ③ 判断 userId === ownerUserId;全部通过才打开 SSE 流。

就算攻击者偷到 taskId,没有用户的 Cookie,访问 SSE 直接 403。 规避了 URL 泄露带来的越权风险。 局限:系统要基于 Cookie 做鉴权,不能用 Authorization Bearer header。

方案 3:taskId + Authorization header(公网,最强安全)

浏览器原生 EventSource 不支持自定义 header,所以不能用 EventSource。 使用 fetch + ReadableStream,手动解析 SSE 协议格式。 流程:

  1. POST 创建任务,Authorization 鉴权,返回 taskId。
  2. fetch GET /sse?taskId=xxx,手动带上 Authorization: Bearer xxx
  3. 后端同时校验:
    • Authorization 解析得到请求者 userId
    • taskId 查询得到任务 ownerUserId
    • 两者相等,且 taskId 未过期,才返回流式响应。

✅ 安全最强,不受 cookie 限制;代价是手写 SSE 解析逻辑,不能用 EventSource。

5. 额外安全细节

  1. 限制 SSE 访问频率,防止暴力枚举 taskId;因为 taskId 是高熵随机串,暴力枚举概率极低,但最好加防护。
  2. CORS:SSE 是 GET 简单请求,如果使用 Cookie 模式,记得开启 withCredentials=true
// EventSource开启携带cookie
const es = new EventSource("/sse?taskId=xxx", {withCredentials:true})

坑点:每次任务就一个 id,用完就删行不行?

绝对不行

每个推理任务一个 taskId,任务结束直接删除 Redis 里这条记录

问题核心:SSE 自动重连发生在任务还没跑完 / 已经跑完两种场景,行为完全不一样

两种重连场景

场景 1:推理还没跑完,网络抖动断开,EventSource 自动重连

用户网断了 2 秒,模型还在输出 token,后端还在跑推理。

  • 浏览器立刻发起新 GET /sse?taskId=xxx
  • ✅ 如果 taskId 还在 Redis:正常恢复 SSE 连接,可以继续推送后续 token。
  • ❗如果你任务没结束就提前删 id:重连过来直接 404/403,流直接断掉,用户拿不到剩下内容。

所以:不能任务一结束输出完才删,更不能中途删;也不能重连一次就销毁 taskId。taskId 是要支持多次连接的,不是一次性 “只能连接一次”。

场景 2:推理已经全部完成,输出全部结束,你删除 taskId

用户已经看完全部回答,SSE 正常 close。 之后用户如果因为某些原因(浏览器自动重试、页面休眠恢复)又触发重连: taskId 已经被删掉 → SSE 直接报错,EventSource 会不断重试,疯狂 403。


关键矛盾

你希望:用完就销毁,防止 taskId 泄露被滥用。 SSE 重连要求:taskId 在任务生命周期内必须可重复访问(允许多次建立 SSE 连接)

不能设计成 “只能连接一次,连上一次就删掉 id”,重连机制直接废掉。

那怎么平衡:允许重连 + 防止泄露滥用?

不要 “连接一次就删”,改成:

  1. taskId 绑定任务状态:pending / running / complete / error,Redis 存,设置全局 TTL 过期(例如 30min),就算忘记删,也会自动失效。
  2. 任务 complete 之后,不要立刻删除 key,只是标记状态为 complete
    • 重连过来依然允许 SSE 连接,把全部历史事件一次性全部推送给客户端,然后正常关闭流。
    • complete 状态保留一小段窗口(比如 5‑10 分钟),供用户断网重连、页面刷新回来看完整结果。
  3. 过了这个保留窗口,再删除 Redis 的 taskId 记录。
  4. 一旦用户主动取消任务,立刻删除 taskId。

简单理解: running 阶段:允许多次 SSE 连接、重连,源源不断推送增量事件。 complete 之后留一个 “回放窗口”:重连就把全部历史一次性推完,然后关闭。窗口到期彻底销毁。

Last‑Event‑Id 和缓存事件的配合(非常重要)

SSE 重连时浏览器会带上 Last‑Event‑Id 请求头,告诉后端:我收到到哪一条事件了,请从这条之后继续发

Redis 存储结构示例:

key: task:{taskId}
{
  userId: 1001,
  status: running | complete | error,
  events: [ 事件1,事件2,事件3,... ], // 已经产出的所有sse事件
  createAt: 1752000000,
  expireAt: 1752001800
}

重连逻辑伪代码:

GET /sse?taskId=xxx
1. 查询task记录,不存在返回403
2. 校验归属权(userId匹配)
3. 读取请求头 Last‑Event‑Id
4. 如果有Last‑Event‑Id:从该id之后取出events,推送给前端
5. 如果status=running:保持长连接,继续推送新产生的事件
6. 如果status=complete:把剩余事件全部推送,直接结束SSE连接
7. 任务complete后,等待5‑10分钟,再删除redis key
  • 如果没有保存 events 数组:重连来了,不知道哪些已经发给客户端,无法断点续传。

代价:你必须缓存该任务已经输出的全部 SSE 事件。

实操建议
  1. ❌不要:SSE 连上就删 taskId;任务一输出完立刻删 taskId。
  2. ✅推荐:complete 标记状态,保留 5‑10 分钟回放窗口,Redis 同时设置 30min 最大 TTL 兜底。窗口结束删除 key。
  3. ✅必须缓存该任务全部 SSE 事件列表,用于重连回放,配合Last‑Event‑Id
  4. ✅公网环境,叠加归属校验:即使 taskId 有效,也要校验当前请求用户是否等于任务 owner。
  5. ✅taskId 使用高熵随机串,禁止自增 ID。
补充一个现实坑

用户打开页面,SSE 跑一半,电脑休眠。唤醒后浏览器 EventSource 会自动重连。 如果此时 taskId 已经被删掉,会不停报 error,疯狂重试,用户看不到剩余内容。保留回放窗口就是为了解决这类场景。

方案二:fetch+ReadableStream+Authorization Bearer

核心优势对比 EventSource:可以自由加任意请求头(Authorization Bearer)、支持 GET/POST、AbortController 可控取消,没有浏览器原生 EventSource 的 header 限制。

每次流式请求带上 Authorization Bearer Token(公网最推荐)

每一次 fetch 请求,都带上 Authorization 头,后端每次都解析 token 拿到 userId,再结合 taskId 做归属校验。

  • 不再依赖 URL 里 taskId 单独做信任凭证;就算 taskId 泄露,攻击者没有合法 Bearer token 也访问不了。
  • 完全避开 EventSource 不能自定义 header 的痛点。

前端关键要点:

  1. AbortController 用来手动终止流(页面卸载、用户点停止)
  2. 跨域场景:携带自定义Authorization会触发 OPTIONS 预检,后端 CORS 必须放行Authorization
  3. fetch 默认不带 cookie,如需 cookie 要加credentials:"include";但 Bearer token 放 header,不需要靠 cookie。
后端鉴权逻辑(重点,双重校验)

GET /api/sse?taskId=xxx

  1. 解析 HTTP Header Authorization: Bearer xxx,校验 token 合法性,得到requestUserId;token 失效直接返回 401。
  2. 根据 taskId 查询 Redis,拿到任务记录:{ownerUserId, status, events, expireAt},taskId 不存在 / 过期返回 403。
  3. 校验归属:requestUserId === ownerUserId,不相等直接 403 拒绝。就算攻击者偷到 taskId,没有合法 Bearer token,第一步就 401;就算拿到 token,task 不属于该用户也 403。安全等级很高。
  4. 校验通过,开启流式 SSE 输出,支持Last‑Event‑Id做断点重放。

注意:每次重连,fetch 会重新完整发送 Authorization 头;不像 EventSource,这里重连鉴权天然生效。

项目 EventSource + taskId fetch + ReadableStream + taskId
自定义 Authorization 头 ❌浏览器原生不支持 ✅完全支持
自动重连 ✅浏览器自带 ❌需要自己写重试、Last‑Event‑Id
取消连接 .close();后端难感知 AbortController.abort();后端能感知 TCP 断开
鉴权方式 只能 Cookie /url 参数 taskId Bearer header / Cookie /taskId 三重自由组合
CORS GET 简单请求,一般无预检 携带 Authorization 触发 OPTIONS 预检
断点续传 原生Last‑Event‑Id 需要前端手动传递Last‑Event‑Id请求头

代码示例

/**
 * fetch + ReadableStream SSE 带自动重连
 * @param {string} taskId 任务id
 * @param {string} token bearer token
 * @param {Object} callbacks {onMessage,onComplete,onError}
 * @returns { {abort:()=>void} }
 */
function createSseStream(taskId, token, callbacks) {
  const { onMessage, onComplete, onError } = callbacks;

  // 状态
  let controller = null;
  let lastEventId = ""; // 记录上一条事件id,用于断点续连
  let retryCount = 0;
  const MAX_RETRY = 5; // 最大重试次数
  let isUserAbort = false; // 用户是否手动取消

  // 指数退避:1s,2s,4s,8s…
  function getBackoffMs(attempt) {
    return Math.min(1000 * (2 ** attempt), 8000);
  }

  async function connect() {
    if (isUserAbort) return;
    controller = new AbortController();

    try {
      const url = `/api/sse?taskId=${encodeURIComponent(taskId)}`;
      const headers = {
        "Accept": "text/event-stream",
        "Authorization": `Bearer ${token}`,
      };
      // 重连的时候带上上次的事件id
      if (lastEventId) {
        headers["Last-Event-Id"] = lastEventId;
      }

      const res = await fetch(url, {
        method: "GET",
        headers,
        signal: controller.signal,
      });

      if (!res.ok) {
        // 401/403/404:权限问题,不要再重试
        throw new Error(`HttpError ${res.status}`);
      }

      // 连接成功,重置重试计数
      retryCount = 0;

      const reader = res.body.getReader();
      const decoder = new TextDecoder("utf-8");
      let buffer = "";

      while (true) {
        const { done, value } = await reader.read();
        if (done) {
          // http流正常结束(服务端complete关闭流)
          onComplete?.();
          return;
        }

        buffer += decoder.decode(value, { stream: true });
        const chunks = buffer.split("\n\n");
        buffer = chunks.pop() || "";

        for (const raw of chunks) {
          if (!raw.trim()) continue;
          // 解析标准SSE单条事件
          const eventObj = parseSseEvent(raw);
          // 更新最后收到的事件id,重连要用
          if (eventObj.id) {
            lastEventId = eventObj.id;
          }

          if (eventObj.event === "complete") {
            onComplete?.();
            return;
          }
          if (eventObj.event === "error") {
            onError?.(new Error(eventObj.data));
            return;
          }
          if (eventObj.data) {
            onMessage(eventObj.data, eventObj.event, eventObj.id);
          }
        }
      }
    } catch (err) {
      if (isUserAbort) {
        // 用户主动abort,不重连
        return;
      }
      onError?.(err);

      // 达到最大重试次数,停止
      if (retryCount >= MAX_RETRY) {
        return;
      }

      retryCount += 1;
      const waitMs = getBackoffMs(retryCount);
      // 等待后再次connect
      setTimeout(() => {
        connect();
      }, waitMs);
    }
  }

  // 解析单条SSE事件,输入是已经按 \n\n 切好的单段
  function parseSseEvent(rawStr) {
    const obj = { event: "message", id: "", data: "" };
    const lines = rawStr.split("\n");
    for (const line of lines) {
      if (line.startsWith("id:")) {
        obj.id = line.slice(3).trimStart();
      } else if (line.startsWith("event:")) {
        obj.event = line.slice(6).trimStart();
      } else if (line.startsWith("data:")) {
        obj.data = line.slice(5);
      }
    }
    return obj;
  }

  // 启动连接
  connect();

  // 返回对外终止接口
  return {
    abort() {
      isUserAbort = true;
      controller?.abort();
    },
  };
}

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐