Vue2 SSE 流式对话完整前端代码
模板部分(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. 前后端交互流程
- POST /api/chat/submit:提交提问,后端缓存对话,返回唯一
sessionId - GET /api/stream/connect?sessionId=xxx:EventSource 建立长连接,浏览器自动携带
Last-Event-ID请求头断点续传 - 后端每条推送携带
id: 数字分片下标,前端自动记录lastChunkId
2. 文本渲染优化(无频繁 DOM 重组)
- 容器
white-space: pre-wrap,后端直接输出\n\n,前端不做换行替换处理 - 每条对话维护独立
fullText完整缓存,增量仅做字符串+=单次追加 - 缓冲区
buffer防抖渲染,遇到###/##/。/!/?强制实时渲染 - Markdown 只在防抖 / 强制节点执行一次渲染,不循环拆分重组 DOM
3. 滚动性能优化
- 滚动置底使用节流 50ms,解决 onmessage 高频触发滚动卡顿
- 容器设置
scroll-behavior: auto,超大文本禁用平滑滚动,减少浏览器重绘开销
4. 多气泡隔离防串文
- 每条 AI 气泡独立
sessionId,sseMap存储多个独立 EventSource 实例 - 分片携带
sessionId+chunkId,数据按会话隔离,多对话不会文字错乱 - 不使用 v-if 销毁气泡,对话列表持久存在,仅更新内部 html 内容
5. 重连与指数退避(解决原生 EventSource 无限快速重试)
- 最大重连 5 次,退避策略:1000ms → 2000ms → 4000ms
- 连接成功重置重连计数;达到上限永久关闭连接,提示网络异常
- 浏览器原生自动携带
Last-Event-ID,后端从对应分片续传,无需重新生成全文
6. 错误分层处理
- 网络断连:进入指数退避重连流程
- 401 鉴权失效:清除 token,跳转登录页,停止重试
- 429 限流 / 额度耗尽:终止重连,展示友好错误提示
- 自定义 event: error 业务推送:立即关闭 SSE,停止接收分片
- 25s 分片超时兜底:调用
/stream/status同步任务状态,卡死场景兜底恢复完整文本
7. 生命周期资源释放(防内存泄漏 / 销毁后回调报错)
beforeDestroy遍历关闭所有 EventSource,手动置空onmessage/onerror/onopen监听- 全局清除滚动、loading、分片超时所有定时器
- 单条会话销毁时同步清除对应重连状态、超时计时器、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 长连接拉流
- 第一步
POST /create-task:提交 prompt、复杂 header、大请求体,后端生成任务,返回taskId - 第二步
GET /sse/stream?taskId=xxx:浏览器用标准EventSource建立 SSE,通过 taskId 订阅这个任务的流式输出
✅ 这个架构确实规避原生 SSE 不能 POST 的最大痛点,是工业界非常经典的折中方案。下面对比:该方案 VS 直接 Streamable (fetch+ReadableStream),讲清收益、坑点、选型。
1.方案优点
- 解决 SSE 只能 GET 的硬伤:复杂入参、大 prompt、鉴权 body 全部放 POST,不走 SSE 的 URL 参数。
- 复用浏览器原生
EventSource:自带自动重连、Last‑Event‑Id、事件类型 event、retry。断网刷新后,SSE 依靠Last‑Event‑ID可以从断点继续接收分片,不用全部重跑大模型推理。 - URL 只有简短 taskId,不会出现 GET 超长 URL、参数泄露在日志、浏览器 URL 长度限制问题。
- 任务和流解耦:任务在后端独立执行,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.方案隐藏坑点
这个方案不是银弹,代价来自任务异步化
- 后端必须缓存输出分片 SSE 断开重连之后,依靠
Last‑Event‑Id恢复,后端要把已经输出的 token / 事件存起来(内存 /redis)。如果不缓存,重连过来拿不到历史分片,断点续传失效。 还要做任务过期清理,防止内存泄漏。 - 任务后台运行,用户关闭 SSE≠停止推理 用户浏览器调用
eventSource.close()只是关闭浏览器侧连接。后端后台推理还在继续跑,消耗 GPU / 算力。 👉 必须额外机制:
- SSE 断开检测;或者前端再发一个 POST
/cancel?taskId=xxx通知后端停止大模型。 否则会大量浪费算力。
- 任务排队、超时状态需要 SSE 推送事件 POST 返回 taskId 的时候,推理可能还没开始、排队中。SSE 通道需要推送自定义事件:
event:pending、event:running、event:error、event:complete。 - 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.成立的前提
- POST 创建任务时:后端校验 Authorization /cookie,确认是谁,校验权限,没问题生成 taskId。
- taskId 和
userId绑定存储(Redis/DB):taskId → {userId, expireAt, status} - SSE GET
/sse?taskId=xxx:- 后端拿到 taskId,查这个任务属于哪个用户,任务是否存在、是否过期。
- 不再解析 Authorization 请求头,依靠 taskId 完成身份映射。
✅ 这样是可行的,很多异步流式系统就是这么干。
⚠️但是:taskId ≠ 随便的业务 ID,它具备凭证属性,有安全风险。
2. 核心安全风险
风险 1:taskId 泄露
taskId 放在URL Query 参数里(SSE 只能放 query):
- 浏览器历史记录保存完整 URL;
- 代理服务器、Nginx 访问日志会打印完整 query 字符串;
- 网络抓包能看到 URL; 如果别人拿到你的 taskId,就可以订阅你的 SSE 流,读取模型输出内容。
普通业务 ID 泄露只是查不到数据;taskId 泄露等价于临时密钥泄露,可以直接拿流式数据。
👉 对策:
- taskId 不能是自增数字,要用高熵随机字符串(uuid/v4 /crypto 随机串),不可猜测,禁止
1,2,3,4这种自增 id。 - 设置短有效期,比如 taskId 5‑30 分钟过期,过期直接作废。
- Redis 存储,设置 TTL 自动过期,不要永久保存。
- 任务完成后立刻作废 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)
- SSE 过来,先拿 Authorization 解析出 userId;
- 根据 taskId 查 redis 得到该任务所属的 ownerUserId;
- 判断:
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 凭证(你设想的)
流程:
- POST /create‑task 带 Authorization → 鉴权成功,生成高熵短时效 taskId,绑定 userId 存入 redis。
- EventSource GET
/sse?taskId=xxx,只传 taskId。 - SSE 后端:查 redis,taskId 是否存在、未过期,存在就允许推送。
优点:实现简单,SSE 不需要处理 token。 缺点:taskId 在 url,日志 / 历史会泄露;一旦泄露直接访问数据。 适用:内网系统、可信环境。公网不推荐直接裸上。
方案 2:taskId + Cookie 双重校验(公网推荐,配合 EventSource)
- POST 创建任务,依靠 Cookie 做用户鉴权,生成 taskId 绑定 userId。
- EventSource 发起 SSE 请求,浏览器自动带上 Cookie。
- 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 协议格式。 流程:
- POST 创建任务,Authorization 鉴权,返回 taskId。
- fetch GET
/sse?taskId=xxx,手动带上Authorization: Bearer xxx。 - 后端同时校验:
- Authorization 解析得到请求者 userId
- taskId 查询得到任务 ownerUserId
- 两者相等,且 taskId 未过期,才返回流式响应。
✅ 安全最强,不受 cookie 限制;代价是手写 SSE 解析逻辑,不能用 EventSource。
5. 额外安全细节
- 限制 SSE 访问频率,防止暴力枚举 taskId;因为 taskId 是高熵随机串,暴力枚举概率极低,但最好加防护。
- 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”,重连机制直接废掉。
那怎么平衡:允许重连 + 防止泄露滥用?
不要 “连接一次就删”,改成:
- taskId 绑定任务状态:
pending / running / complete / error,Redis 存,设置全局 TTL 过期(例如 30min),就算忘记删,也会自动失效。 - 任务 complete 之后,不要立刻删除 key,只是标记状态为 complete。
- 重连过来依然允许 SSE 连接,把全部历史事件一次性全部推送给客户端,然后正常关闭流。
- complete 状态保留一小段窗口(比如 5‑10 分钟),供用户断网重连、页面刷新回来看完整结果。
- 过了这个保留窗口,再删除 Redis 的 taskId 记录。
- 一旦用户主动取消任务,立刻删除 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 事件。
实操建议
- ❌不要:SSE 连上就删 taskId;任务一输出完立刻删 taskId。
- ✅推荐:complete 标记状态,保留 5‑10 分钟回放窗口,Redis 同时设置 30min 最大 TTL 兜底。窗口结束删除 key。
- ✅必须缓存该任务全部 SSE 事件列表,用于重连回放,配合
Last‑Event‑Id。 - ✅公网环境,叠加归属校验:即使 taskId 有效,也要校验当前请求用户是否等于任务 owner。
- ✅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 的痛点。
前端关键要点:
AbortController用来手动终止流(页面卸载、用户点停止)- 跨域场景:携带自定义
Authorization会触发 OPTIONS 预检,后端 CORS 必须放行Authorization头 - fetch 默认不带 cookie,如需 cookie 要加
credentials:"include";但 Bearer token 放 header,不需要靠 cookie。
后端鉴权逻辑(重点,双重校验)
GET /api/sse?taskId=xxx
- 解析 HTTP Header
Authorization: Bearer xxx,校验 token 合法性,得到requestUserId;token 失效直接返回 401。 - 根据 taskId 查询 Redis,拿到任务记录:
{ownerUserId, status, events, expireAt},taskId 不存在 / 过期返回 403。 - 校验归属:
requestUserId === ownerUserId,不相等直接 403 拒绝。就算攻击者偷到 taskId,没有合法 Bearer token,第一步就 401;就算拿到 token,task 不属于该用户也 403。安全等级很高。 - 校验通过,开启流式 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();
},
};
}
更多推荐




所有评论(0)