Vue 基础体系 · 第 48/70 篇。示例基于 Vue 3、Composition API、TypeScript 与现代 Vite 工具链;版本敏感能力会单独标注。
Vue 实时数据:WebSocket、SSE、重连、心跳和状态同步
实时数据不是“建立一条长连接,然后在 onmessage 中改几个变量”这么简单。真正的实时系统至少要回答五个问题:
- 客户端如何与服务端建立双向或单向连接?
- 连接断开后,什么时候重连,重连多少次,如何避免连接风暴?
- 如何判断连接只是暂时空闲,还是已经不可用?
- 重连期间丢失的数据如何恢复?
- 多条消息并发到达、重复到达或乱序到达时,Vue 中的状态如何保持正确?
本文基于 Vue 3、Composition API、TypeScript 和现代 Vite 工具链,系统说明 WebSocket、SSE、重连、心跳与状态同步之间的关系。
一、先区分“实时”与“传输方式”
“实时”不是某个协议名称,而是一种业务要求:服务端状态变化后,客户端应在可接受的延迟内得到通知。
常见实现有三类:
| 方式 | 方向 | 浏览器 API | 适合场景 |
|---|---|---|---|
| 轮询 | 客户端定期请求 | fetch、Axios |
实现简单、更新不频繁 |
| SSE | 服务端到客户端 | EventSource |
通知、日志、进度、行情推送 |
| WebSocket | 客户端与服务端双向 | WebSocket |
聊天、协同编辑、实时控制 |
这里的“方向”指应用消息的主要流向:
- SSE 的连接由客户端发起,但建立后主要是服务端持续发送数据。
- WebSocket 建立后,客户端和服务端都可以主动发送消息。
- WebSocket 握手通常从 HTTP 请求开始,随后升级为独立的双向连接。
- SSE 仍然是 HTTP 响应流,服务端以
text/event-stream格式持续输出事件。
因此,选择协议时不能只看“哪个更实时”,而应看是否需要客户端主动向服务端发送实时业务消息。
二、WebSocket 的连接模型
2.1 WebSocket 的生命周期
浏览器原生 WebSocket 对象有四种标准状态:
WebSocket.CONNECTING // 0,正在连接
WebSocket.OPEN // 1,连接已建立
WebSocket.CLOSING // 2,正在关闭
WebSocket.CLOSED // 3,已经关闭
典型生命周期如下:
stateDiagram-v2
[*] --> CONNECTING
CONNECTING --> OPEN: open
CONNECTING --> CLOSED: error/close
OPEN --> CLOSING: close()
OPEN --> CLOSED: close/error
CLOSING --> CLOSED: close
CLOSED --> CONNECTING: 重新创建 WebSocket
一个重要事实是:关闭后的 WebSocket 对象不能重新打开。重连必须创建新的实例:
socket = new WebSocket(url)
不能这样做:
socket.close()
// 错误思路:试图让同一个 socket 再次 open
socket.onopen = ...
2.2 open、message、error 和 close
浏览器 WebSocket 的几个事件职责不同:
open:TCP/TLS 和 WebSocket 握手成功,可以发送应用消息。message:收到服务端数据。error:发生通信错误,但错误对象通常不包含足够的诊断细节。close:连接最终关闭,包含关闭码和原因。
实际重连通常应放在 close 中,而不是只放在 error 中,因为错误后通常还会触发关闭流程,而且有些关闭场景不会先提供有用的 error 信息。
socket.onerror = (event) => {
console.warn('WebSocket error', event)
}
socket.onclose = (event) => {
console.log('WebSocket closed', {
code: event.code,
reason: event.reason,
wasClean: event.wasClean,
})
// 在这里决定是否重连
}
event.wasClean 只能说明关闭过程是否符合协议和实现预期,不能直接等价于“业务状态没有问题”。例如服务端正常关闭连接,业务上的订阅状态仍然可能需要重新恢复。
2.3 WebSocket 不等于业务协议
WebSocket 只定义了连接和帧传输机制,不会自动定义以下内容:
- 消息是 JSON、文本还是二进制;
- 消息类型如何区分;
- 服务端推送的版本号是什么;
- 客户端如何确认收到消息;
- 断线后如何补发数据;
- 心跳内容是什么。
因此,应用通常要定义一层业务协议:
type ServerMessage =
| {
type: 'snapshot'
stream: string
version: number
data: Order[]
}
| {
type: 'order.updated'
stream: string
version: number
data: Order
}
| {
type: 'pong'
requestId: string
}
| {
type: 'error'
code: string
message: string
}
其中:
type表示消息类型;stream表示消息属于哪个数据流;version表示服务端状态版本;data是快照或增量数据;requestId用于把响应与请求关联起来。
如果没有消息类型和版本号,客户端很难判断一条消息应该如何应用。
三、SSE:服务端到客户端的事件流
3.1 SSE 的数据方向
SSE 的全称是 Server-Sent Events。它使用 HTTP 长响应,把服务端事件按照特定文本格式持续发送给浏览器。
服务端响应通常需要包含:
Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
最小事件格式如下:
event: order.updated
id: 42
data: {"id":"o-1","status":"paid"}
事件末尾的空行表示一条事件结束。data 可以出现多次,多行内容会被拼接。
客户端使用:
const source = new EventSource('/api/events')
SSE 的主要特点是:
- 浏览器原生支持自动重连;
- 服务端可以发送命名事件;
- 客户端通过
message或addEventListener接收消息; - 客户端不能通过同一个
EventSource连接发送业务消息; - 若客户端需要发送命令,仍需使用
fetch、普通 HTTP 请求或另一个 WebSocket。
3.2 SSE 的 id 与 Last-Event-ID
SSE 事件可以携带 id:
id: 42
data: {"status":"paid"}
浏览器会记录最近收到的事件 ID。连接意外断开后,浏览器通常会在重连请求中发送:
Last-Event-ID: 42
服务端可以根据这个 ID 补发后续事件。
这提供了“从上次位置继续”的基础,但它不是完整的数据一致性保证。服务端仍然必须:
- 保存一段时间的事件历史;
- 判断客户端的 ID 是否仍在可补发范围内;
- 超出范围时发送完整快照;
- 保证事件 ID 在同一数据流中单调递增。
如果服务端只实时广播、不保存历史,那么 Last-Event-ID 没有补偿作用。
3.3 SSE 的客户端示例
import { onBeforeUnmount, ref } from 'vue'
interface Order {
id: string
status: string
}
export function useOrderEvents() {
const orders = ref<Order[]>([])
const connected = ref(false)
const error = ref<Event | null>(null)
const source = new EventSource('/api/orders/events', {
withCredentials: true,
})
source.onopen = () => {
connected.value = true
error.value = null
}
source.onerror = (event) => {
connected.value = false
error.value = event
// EventSource 默认会尝试重连。
// 这里通常只更新 UI 或记录诊断信息。
}
source.addEventListener('order.snapshot', (event) => {
const message = JSON.parse((event as MessageEvent).data) as {
version: number
data: Order[]
}
orders.value = message.data
})
source.addEventListener('order.updated', (event) => {
const message = JSON.parse((event as MessageEvent).data) as {
data: Order
}
const index = orders.value.findIndex(
(order) => order.id === message.data.id,
)
if (index === -1) {
orders.value.push(message.data)
} else {
orders.value[index] = message.data
}
})
onBeforeUnmount(() => {
source.close()
})
return {
orders,
connected,
error,
}
}
这里使用 source.close() 是必要的。否则组件卸载后,连接可能继续存在,导致:
- 页面已经没有组件消费消息;
- 仍然占用浏览器和服务端连接;
- 旧组件继续修改状态;
- 路由切换后产生重复订阅。
3.4 SSE 的限制
SSE 不是 WebSocket 的“更简单版本”,它的边界不同:
- 浏览器原生
EventSource不支持自定义请求头; - 认证通常使用同源 Cookie,或把短期凭证放到 URL,但 URL 凭证容易出现在日志中;
- 需要自定义请求头时,可以使用基于
fetch的事件流实现,但那不是原生EventSource; - 代理、负载均衡器和压缩中间件必须允许响应持续刷新;
- 服务端应定期发送注释或事件,避免中间代理误判连接空闲。
SSE 的注释心跳可以是:
: keep-alive
这不是业务事件,但会刷新 HTTP 流。它能帮助保持连接活跃,却不能自动证明业务逻辑正常工作。
四、WebSocket 与 SSE 如何选择
可以用一个直接的判断过程:
- 如果只有服务端向客户端推送,先考虑 SSE。
- 如果客户端也需要低延迟发送消息,考虑 WebSocket。
- 如果需要浏览器原生自动重连,并且消息是事件流,SSE 更简单。
- 如果需要二进制帧、双向请求响应、订阅和取消订阅,WebSocket 更合适。
- 如果部署环境对长 HTTP 响应支持不稳定,SSE 需要重点验证代理和网关配置。
- 如果客户端需要带复杂认证头,原生
EventSource可能不合适。
二者都不能自动解决状态同步问题。协议只负责传输,数据是否完整、是否重复、是否乱序,需要应用协议负责。
五、重连:从“重新连接”到可控算法
5.1 为什么不能立即无限重连
假设服务端宕机,客户端在 close 后立即重连:
function reconnect() {
socket = new WebSocket(url)
}
如果每个客户端都这样做,结果会是:
- 所有连接几乎同时失败;
- 所有客户端几乎同时发起新连接;
- 服务端仍未恢复,连接再次失败;
- 客户端进入高频重试;
- 服务端恢复时遭遇连接洪峰。
这种现象称为连接风暴。解决它需要退避和抖动。
5.2 指数退避与抖动
设第 次重连等待时间为:
其中:
- 是初始延迟,例如 1000 毫秒;
- 是最大延迟,例如 30 秒;
- 是已经失败的重连次数;
- 是随机抖动;
- 常见全抖动方式为 的随机值。
完整算法可以写成:
function getReconnectDelay(
attempt: number,
base = 1000,
max = 30_000,
): number {
const exponential = Math.min(max, base * 2 ** attempt)
const jitter = Math.random() * exponential
return exponential + jitter
}
假设 base = 1000、max = 30000:
| 失败次数 | 指数部分 | 实际等待范围 |
|---|---|---|
| 0 | 1 秒 | 1~2 秒 |
| 1 | 2 秒 | 2~4 秒 |
| 2 | 4 秒 | 4~8 秒 |
| 3 | 8 秒 | 8~16 秒 |
| 4 | 16 秒 | 16~32 秒 |
| 5 | 30 秒 | 30~60 秒 |
实际项目可以把随机范围限制在固定比例内,但必须让不同客户端的重连时间错开。
5.3 哪些关闭不应该重连
不是所有关闭都代表临时故障:
- 用户主动退出登录;
- 页面或模块主动销毁连接;
- 服务端返回“认证失效”;
- 服务端明确要求客户端升级协议;
- 当前页面已经离开,不再需要该数据流。
因此需要区分“主动关闭”和“意外关闭”:
let manuallyClosed = false
function stop() {
manuallyClosed = true
socket?.close(1000, 'client shutdown')
}
function handleClose() {
if (manuallyClosed) {
return
}
scheduleReconnect()
}
如果不做这个区分,组件卸载时调用 close(),onclose 又启动重连,就会出现“页面离开后连接仍然不断创建”的错误。
六、一个可复用的 WebSocket Composition API
下面的实现包含:
- 连接状态;
- 指数退避和随机抖动;
-应用层心跳; - 主动关闭与意外关闭区分;
- 连接代次,避免旧连接事件干扰新连接;
- JSON 消息发送和接收。
import {
onBeforeUnmount,
readonly,
ref,
type Ref,
} from 'vue'
export type RealtimeStatus =
| 'idle'
| 'connecting'
| 'open'
| 'closed'
| 'error'
interface UseWebSocketOptions<T> {
url: string
heartbeatInterval?: number
reconnect?: boolean
reconnectBaseDelay?: number
reconnectMaxDelay?: number
onMessage?: (message: T) => void
}
export function useWebSocket<TSend extends object, TReceive>(
options: UseWebSocketOptions<TReceive>,
) {
const status = ref<RealtimeStatus>('idle')
const lastError = ref<Event | null>(null)
const reconnectAttempt = ref(0)
let socket: WebSocket | null = null
let reconnectTimer: ReturnType<typeof setTimeout> | undefined
let heartbeatTimer: ReturnType<typeof setInterval> | undefined
let stopped = false
let generation = 0
const heartbeatInterval = options.heartbeatInterval ?? 20_000
const reconnectBaseDelay = options.reconnectBaseDelay ?? 1_000
const reconnectMaxDelay = options.reconnectMaxDelay ?? 30_000
function clearReconnectTimer() {
if (reconnectTimer !== undefined) {
clearTimeout(reconnectTimer)
reconnectTimer = undefined
}
}
function clearHeartbeat() {
if (heartbeatTimer !== undefined) {
clearInterval(heartbeatTimer)
heartbeatTimer = undefined
}
}
function reconnectDelay(attempt: number) {
const exponential = Math.min(
reconnectMaxDelay,
reconnectBaseDelay * 2 ** attempt,
)
return exponential + Math.random() * exponential
}
function startHeartbeat(current: WebSocket, currentGeneration: number) {
clearHeartbeat()
heartbeatTimer = setInterval(() => {
if (
currentGeneration !== generation ||
current !== socket ||
current.readyState !== WebSocket.OPEN
) {
return
}
// 这是应用层 ping,不是浏览器原生 WebSocket 协议 ping。
current.send(JSON.stringify({ type: 'ping' }))
}, heartbeatInterval)
}
function scheduleReconnect() {
if (stopped || options.reconnect === false) {
return
}
clearReconnectTimer()
const attempt = reconnectAttempt.value
const delay = reconnectDelay(attempt)
reconnectTimer = setTimeout(() => {
reconnectTimer = undefined
reconnectAttempt.value += 1
connect()
}, delay)
}
function connect() {
if (stopped) {
return
}
clearReconnectTimer()
clearHeartbeat()
const currentGeneration = ++generation
const current = new WebSocket(options.url)
socket = current
status.value = 'connecting'
current.onopen = () => {
// 旧连接的延迟事件不能修改当前连接状态。
if (
currentGeneration !== generation ||
current !== socket
) {
return
}
status.value = 'open'
lastError.value = null
reconnectAttempt.value = 0
startHeartbeat(current, currentGeneration)
}
current.onmessage = (event) => {
if (
currentGeneration !== generation ||
current !== socket
) {
return
}
try {
const message = JSON.parse(event.data) as TReceive
options.onMessage?.(message)
} catch (error) {
console.error('Invalid WebSocket message', {
eventData: event.data,
error,
})
}
}
current.onerror = (event) => {
if (
currentGeneration !== generation ||
current !== socket
) {
return
}
status.value = 'error'
lastError.value = event
}
current.onclose = () => {
if (
currentGeneration !== generation ||
current !== socket
) {
return
}
clearHeartbeat()
status.value = 'closed'
scheduleReconnect()
}
}
function send(message: TSend): boolean {
if (!socket || socket.readyState !== WebSocket.OPEN) {
return false
}
socket.send(JSON.stringify(message))
return true
}
function stop() {
stopped = true
clearReconnectTimer()
clearHeartbeat()
// 让旧连接的事件失效。
generation += 1
const current = socket
socket = null
if (
current &&
(current.readyState === WebSocket.OPEN ||
current.readyState === WebSocket.CONNECTING)
) {
current.close(1000, 'client shutdown')
}
status.value = 'closed'
}
function start() {
if (!stopped) {
return
}
stopped = false
connect()
}
onBeforeUnmount(stop)
connect()
return {
status: readonly(status) as Readonly<Ref<RealtimeStatus>>,
lastError: readonly(lastError),
reconnectAttempt: readonly(reconnectAttempt),
send,
start,
stop,
}
}
6.1 为什么需要 generation
假设发生以下顺序:
- 连接 A 断开;
- 客户端创建连接 B;
- 连接 A 的某个延迟事件才到达;
- 连接 A 的回调修改了共享状态。
如果没有代次检查,旧连接可能把当前状态错误地改成 closed,或者把旧消息写入新连接的数据流。
generation 的作用是为每次连接分配代号。只有当前代连接的事件可以修改状态:
if (currentGeneration !== generation) {
return
}
这不是 WebSocket 规范要求,而是异步 JavaScript 中防止旧任务污染新状态的工程机制。
6.2 在组件中使用
<script setup lang="ts">
import { computed, ref } from 'vue'
import {
useWebSocket,
type RealtimeStatus,
} from '@/composables/useWebSocket'
interface Order {
id: string
status: 'pending' | 'paid' | 'cancelled'
}
type ServerMessage =
| {
type: 'order.updated'
version: number
data: Order
}
| {
type: 'pong'
}
const orders = ref<Order[]>([])
const version = ref(0)
const {
status,
lastError,
send,
stop,
start,
} = useWebSocket<{
type: 'subscribe'
stream: string
}, ServerMessage>({
url: 'wss://example.com/realtime',
onMessage(message) {
if (message.type === 'order.updated') {
if (message.version <= version.value) {
return
}
const index = orders.value.findIndex(
(order) => order.id === message.data.id,
)
if (index === -1) {
orders.value.push(message.data)
} else {
orders.value[index] = message.data
}
version.value = message.version
}
},
})
const statusText = computed(() => {
const map: Record<RealtimeStatus, string> = {
idle: '未连接',
connecting: '连接中',
open: '已连接',
closed: '已关闭',
error: '连接错误',
}
return map[status.value]
})
function subscribe() {
send({
type: 'subscribe',
stream: 'orders',
})
}
</script>
<template>
<section>
<p>连接状态:{{ statusText }}</p>
<p v-if="lastError">连接出现错误,请查看网络诊断。</p>
<button @click="subscribe">订阅订单</button>
<button @click="stop">停止连接</button>
<button @click="start">重新连接</button>
<ul>
<li v-for="order in orders" :key="order.id">
{{ order.id }}:{{ order.status }}
</li>
</ul>
</section>
</template>
这里 send() 返回布尔值:
true:消息已提交给浏览器 WebSocket;false:当前还不是OPEN状态。
但 true 不等于服务端已经处理成功。若业务需要可靠命令,还应定义 requestId、服务端确认消息和超时处理。
七、心跳:检测连接是否“真正可用”
7.1 TCP 连接存在不代表应用可用
一个连接可能处于以下状态:
- 浏览器认为连接仍然是
OPEN; - 网络路径实际已经断开;
- 中间代理丢弃了空闲连接;
- 服务端线程阻塞,无法处理消息;
- 页面进入后台,定时器被浏览器节流。
如果长时间没有业务消息,客户端无法仅凭“没有收到消息”判断连接是否正常,因为业务本来可能就没有更新。
心跳的目标是检测:
客户端发送探测后,服务端是否能在规定时间内返回确认。
7.2 WebSocket 的协议级 Ping 与浏览器限制
WebSocket 协议本身定义了 Ping/Pong 控制帧,但浏览器原生 WebSocket API 通常不向 JavaScript 暴露发送协议级 Ping 的方法。
因此,浏览器前端通常发送应用层消息:
{"type":"ping","requestId":"abc-123"}
服务端返回:
{"type":"pong","requestId":"abc-123"}
应用层 ping 与协议级 Ping 不是同一件事:
- 协议级 Ping 由 WebSocket 实现处理;
- 应用层 Ping 是普通业务数据帧;
- 服务端必须显式解析应用层 Ping 并返回 Pong。
7.3 心跳必须有超时,不只是定时发送
只发送 Ping 而不检查 Pong,无法发现半开连接。完整逻辑应包括:
- 每隔 毫秒发送 Ping;
- 记录本次 Ping 的时间和请求 ID;
- 在 毫秒内等待 Pong;
- 超时后主动关闭连接;
- 由关闭流程触发重连。
可以定义:
例如:
- 心跳间隔 秒;
- Pong 超时 秒。
若客户端连续发送 Ping,但服务端都没有响应,超过 10 秒后关闭连接,比等待浏览器底层网络错误更快恢复。
一个带超时的心跳实现如下:
let heartbeatTimer: ReturnType<typeof setInterval> | undefined
let heartbeatTimeout: ReturnType<typeof setTimeout> | undefined
function startHeartbeat(socket: WebSocket) {
stopHeartbeat()
heartbeatTimer = setInterval(() => {
if (socket.readyState !== WebSocket.OPEN) {
return
}
const requestId = crypto.randomUUID()
socket.send(JSON.stringify({
type: 'ping',
requestId,
}))
heartbeatTimeout = setTimeout(() => {
// 没有收到对应 Pong,主动断开。
// onclose 中统一执行重连。
if (socket.readyState === WebSocket.OPEN) {
socket.close(4000, 'heartbeat timeout')
}
}, 10_000)
}, 20_000)
}
function handlePong(requestId: string) {
// 生产代码应只清除对应 requestId 的超时。
// 若同时允许多个探测,应使用 Map<string, timer> 管理。
clearTimeout(heartbeatTimeout)
heartbeatTimeout = undefined
console.debug('pong received', requestId)
}
function stopHeartbeat() {
if (heartbeatTimer !== undefined) {
clearInterval(heartbeatTimer)
heartbeatTimer = undefined
}
if (heartbeatTimeout !== undefined) {
clearTimeout(heartbeatTimeout)
heartbeatTimeout = undefined
}
}
上面代码为了展示机制,简化了多个并发 Ping 的管理。更严格的实现应使用:
const pendingPings = new Map<
string,
ReturnType<typeof setTimeout>
>()
收到 Pong 时只清除对应的 requestId,避免较早 Ping 的超时任务误关闭连接。
7.4 心跳间隔的取舍
心跳太频繁会:
- 增加服务端连接处理开销;
- 增加移动端耗电;
- 增加网络流量。
心跳太稀疏则会:
- 延迟发现断线;
- 让用户更久看到过期数据;
- 使代理的空闲超时更容易触发。
心跳间隔应结合:
- 网关空闲超时时间;
- 业务允许的最大数据陈旧时间;
- 客户端数量;
- 移动端和后台页面行为。
后台标签页中的定时器可能被浏览器延迟执行,所以前端心跳只能提供尽力而为的检测,不能当作精确计时器。
八、重连后最关键的问题:状态同步
重连只能恢复“连接”,不能自动恢复“连接断开期间发生的事件”。
假设服务端状态版本依次为:
100 -> 101 -> 102 -> 103
客户端已收到版本 100,随后断线。重连后直接收到版本 103。如果 101 和 102 是两个独立业务变化,客户端不能仅凭 103 推断中间发生了什么。
因此,实时系统通常采用以下两种模型:
8.1 快照模型
服务端在连接建立或订阅成功后发送完整状态:
{
"type": "snapshot",
"stream": "orders",
"version": 103,
"data": [
{"id": "o-1", "status": "paid"},
{"id": "o-2", "status": "pending"}
]
}
客户端收到快照后直接替换本地状态。
优点:
- 实现简单;
- 重连恢复容易;
- 不需要客户端保存复杂的事件历史。
缺点:
- 数据量大;
- 快照生成和传输成本高;
- 高频更新时可能浪费带宽。
8.2 增量事件模型
服务端发送每次变化:
{
"type": "order.updated",
"stream": "orders",
"version": 101,
"data": {
"id": "o-1",
"status": "paid"
}
}
客户端按版本应用事件:
function applyEvent(message: {
version: number
data: Order
}) {
if (message.version <= currentVersion) {
// 重复消息或旧消息,忽略
return
}
if (message.version !== currentVersion + 1) {
// 发现缺口,不能继续假设本地状态正确
requestResync()
return
}
updateOrder(message.data)
currentVersion = message.version
}
这里使用了三个条件:
version <= currentVersion:旧消息或重复消息;version === currentVersion + 1:正好是下一条;version > currentVersion + 1:存在缺失事件,需要补偿。
8.3 为什么版本号不能只用客户端时间
客户端时间不适合充当全局顺序:
- 不同机器时钟可能不一致;
- 网络延迟会导致到达顺序与生成顺序不同;
- 多个服务实例的本地时间不能天然提供严格顺序;
- 时间戳相同会造成无法比较。
版本号应由服务端数据源生成,或使用服务端认可的单调序列。若系统是多分区的,则应使用“分区 + 分区内序列”,不能假设不同分区之间有天然全序。
九、快照与增量事件的完整恢复流程
一个可靠的数据流可以采用如下流程:
sequenceDiagram
participant C as Vue 客户端
participant S as 实时服务
participant DB as 数据存储
C->>S: 订阅 stream=orders, lastVersion=100
S->>DB: 查询版本 101 之后的事件
DB-->>S: 101, 102, 103
S-->>C: event 101
S-->>C: event 102
S-->>C: event 103
C->>C: version=103
Note over C,S: 连接断开,期间产生 104、105
C->>S: 重连并订阅 lastVersion=103
S->>DB: 查询版本 104 之后的事件
DB-->>S: 104, 105
S-->>C: event 104
S-->>C: event 105
C->>C: version=105
如果服务端无法找到 103 之后的完整事件历史:
sequenceDiagram
participant C as Vue 客户端
participant S as 实时服务
C->>S: 订阅 lastVersion=103
S-->>C: resync_required
C->>S: 请求 snapshot
S-->>C: snapshot version=205
C->>C: 替换本地状态,version=205
关键原则是:
客户端只有在确认版本连续,或成功应用一个权威快照后,才能认为本地状态可用。
9.1 客户端状态机
可以把数据流状态表示为:
type SyncStatus =
| 'idle'
| 'connecting'
| 'syncing'
| 'live'
| 'stale'
| 'error'
状态转移示例:
idle
-> connecting 开始建立连接
connecting
-> syncing 连接成功,等待快照或补偿事件
syncing
-> live 快照完成或增量连续应用完成
live
-> stale 连接断开或发现版本缺口
stale
-> syncing 重连并开始恢复
syncing
-> error 恢复失败或认证失败
“连接已打开”与“数据已同步”是两个不同状态。UI 不应仅根据 WebSocket.OPEN 显示“数据最新”。
十、在 Pinia 中集中管理实时状态
如果只有一个组件使用数据,可以用 Composition API 局部管理。如果多个页面需要共享同一个实时数据流,应把状态放入 Pinia,避免每个组件各自建立连接。
// stores/orders.ts
import { defineStore } from 'pinia'
import { computed, ref } from 'vue'
import { useWebSocket } from '@/composables/useWebSocket'
interface Order {
id: string
status: 'pending' | 'paid' | 'cancelled'
}
type ServerMessage =
| {
type: 'snapshot'
version: number
data: Order[]
}
| {
type: 'order.updated'
version: number
data: Order
}
| {
type: 'resync_required'
}
| {
type: 'pong'
requestId: string
}
export const useOrdersStore = defineStore('orders', () => {
const orders = ref<Order[]>([])
const version = ref(0)
const syncStatus = ref<
'idle' | 'syncing' | 'live' | 'stale' | 'error'
>('idle')
const byId = computed(() => {
return new Map(orders.value.map((order) => [order.id, order]))
})
const realtime = useWebSocket<
{ type: 'subscribe'; stream: string; lastVersion: number },
ServerMessage
>({
url: 'wss://example.com/realtime',
onMessage(message) {
if (message.type === 'snapshot') {
orders.value = message.data
version.value = message.version
syncStatus.value = 'live'
return
}
if (message.type === 'order.updated') {
if (message.version <= version.value) {
return
}
if (message.version !== version.value + 1) {
syncStatus.value = 'stale'
requestSnapshot()
return
}
const index = orders.value.findIndex(
(order) => order.id === message.data.id,
)
if (index === -1) {
orders.value.push(message.data)
} else {
orders.value[index] = message.data
}
version.value = message.version
syncStatus.value = 'live'
return
}
if (message.type === 'resync_required') {
syncStatus.value = 'syncing'
requestSnapshot()
}
},
})
function subscribe() {
syncStatus.value = 'syncing'
realtime.send({
type: 'subscribe',
stream: 'orders',
lastVersion: version.value,
})
}
function requestSnapshot() {
realtime.send({
type: 'subscribe',
stream: 'orders',
lastVersion: 0,
})
}
function stop() {
realtime.stop()
syncStatus.value = 'stale'
}
return {
orders,
byId,
version,
syncStatus,
subscribe,
requestSnapshot,
stop,
connectionStatus: realtime.status,
}
})
这个示例表达了一个重要的架构关系:
WebSocket/SSE
↓
消息解析与协议校验
↓
版本检查与去重
↓
Pinia store
↓
Vue 组件
组件不应直接承担“连接、重连、补偿、去重、快照替换”等所有职责,否则多个页面很容易建立重复连接,且各自拥有不一致的数据副本。
需要注意,Pinia store 的创建时机与应用生命周期有关。若在服务端渲染环境中使用,连接不能在服务端渲染阶段建立;应在客户端生命周期中启动,并确保不同用户请求之间不共享服务端状态。
十一、Vue 生命周期与路由范围
实时连接到底属于组件、路由页面,还是整个应用,取决于数据流范围。
11.1 页面级连接
如果数据只在某个页面使用,可以在页面组件或组合式函数中创建:
import { onMounted, onBeforeUnmount } from 'vue'
let stopRealtime: (() => void) | undefined
onMounted(() => {
stopRealtime = startRealtime()
})
onBeforeUnmount(() => {
stopRealtime?.()
})
页面离开后关闭连接,资源边界清晰。
11.2 应用级连接
如果顶部通知、消息未读数和多个页面都依赖同一数据流,应由 Pinia store 或应用级服务统一维护。路由切换时不应反复断开和连接。
Vue Router 只负责路由状态和导航生命周期,不会自动管理 WebSocket 或 SSE。可以在路由守卫中决定“是否需要某个订阅”,但连接本身仍需显式创建和销毁。
例如:
router.afterEach((to) => {
const ordersStore = useOrdersStore()
if (to.meta.requiresOrdersRealtime) {
ordersStore.subscribe()
}
})
这段逻辑仍需防止重复订阅。subscribe() 最好具备幂等性,或者记录当前订阅集合:
const subscribedStreams = new Set<string>()
function subscribeOnce(stream: string) {
if (subscribedStreams.has(stream)) {
return
}
subscribedStreams.add(stream)
sendSubscribe(stream)
}
11.3 KeepAlive 的边界
组件被 <KeepAlive> 缓存时,onBeforeUnmount 不一定在离开当前视图时执行。此时应根据业务决定:
- 使用
onActivated恢复订阅; - 使用
onDeactivated暂停订阅; - 或把连接提升到 Pinia,完全脱离页面组件生命周期。
如果忽略这一点,缓存页面可能继续接收实时消息,也可能在重新激活时重复注册事件处理器。
十二、并发、重复与乱序消息
12.1 重复消息
重复消息可能来自:
- 服务端重试发送;
- 客户端重连后补发;
- 网络层或业务层采用至少一次投递;
- 多个订阅重复建立。
因此,更新逻辑不能简单写成:
orders.value.push(message.data)
应根据业务主键和版本去重:
function applyOrder(order: Order, version: number) {
if (version <= currentVersion) {
return
}
const old = orderMap.get(order.id)
if (!old || version >= old.version) {
orderMap.set(order.id, {
data: order,
version,
})
}
currentVersion = Math.max(currentVersion, version)
}
12.2 乱序消息
如果消息 A 的版本是 10,消息 B 的版本是 11,但 B 先到达,直接应用 B 会造成中间状态不确定。
有三种处理策略:
- 严格连续策略:发现版本不是
current + 1就暂停并重新同步。 - 客户端缓冲策略:暂存高版本消息,等待缺失版本到达。
- 按实体版本策略:每个实体独立比较版本,只应用该实体更新版本更高的消息。
选择取决于业务:
- 金融余额、库存、权限等状态通常需要严格连续或服务端快照;
- 独立商品卡片可以按实体版本更新;
- 聊天消息可能需要全局序列保证显示顺序。
12.3 本地操作与服务端推送并发
假设用户点击“支付”:
- 客户端立即把订单显示为
paid; - 服务端稍后推送版本 50,状态仍是
pending; - 客户端如果无条件覆盖,就会出现状态回退。
常见解决方案有:
- 不做乐观更新,等待服务端确认;
- 为本地命令分配
requestId,服务端确认后再提交; - 为每条实体维护服务端版本;
- 将本地待确认操作单独存放,不直接当作权威状态。
一个请求响应协议可以是:
{
"type": "command",
"requestId": "cmd-123",
"name": "pay-order",
"payload": {"orderId": "o-1"}
}
服务端返回:
{
"type": "command.accepted",
"requestId": "cmd-123",
"version": 51
}
或:
{
"type": "command.rejected",
"requestId": "cmd-123",
"reason": "order already cancelled"
}
requestId 还可以帮助服务端实现幂等:同一个命令因重试被发送多次时,服务端只执行一次。
十三、一个最小的 WebSocket 服务端协议示例
以下示例使用 Node.js 与 ws 包说明协议形态。它不是完整生产服务,但可以帮助验证客户端的连接、Ping/Pong 和订阅流程。
安装:
npm install ws
npm install -D tsx @types/ws
server.ts:
import { WebSocketServer, WebSocket } from 'ws'
const wss = new WebSocketServer({ port: 8080 })
let version = 0
const orders = new Map<string, {
id: string
status: 'pending' | 'paid'
}>()
orders.set('o-1', { id: 'o-1', status: 'pending' })
function send(ws: WebSocket, message: unknown) {
if (ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify(message))
}
}
wss.on('connection', (ws) => {
ws.on('message', (raw) => {
let message: unknown
try {
message = JSON.parse(raw.toString())
} catch {
send(ws, {
type: 'error',
code: 'INVALID_JSON',
message: 'message must be JSON',
})
return
}
if (
typeof message !== 'object' ||
message === null ||
!('type' in message)
) {
send(ws, {
type: 'error',
code: 'INVALID_MESSAGE',
message: 'missing message type',
})
return
}
const type = (message as { type: string }).type
if (type === 'ping') {
send(ws, {
type: 'pong',
requestId:
'requestId' in message
? (message as { requestId?: string }).requestId
: undefined,
})
return
}
if (type === 'subscribe') {
send(ws, {
type: 'snapshot',
stream: 'orders',
version,
data: [...orders.values()],
})
return
}
send(ws, {
type: 'error',
code: 'UNKNOWN_TYPE',
message: `unsupported type: ${type}`,
})
})
})
// 每 5 秒模拟一次服务端状态变化。
setInterval(() => {
const order = orders.get('o-1')
if (!order) return
order.status = order.status === 'pending' ? 'paid' : 'pending'
version += 1
const message = JSON.stringify({
type: 'order.updated',
stream: 'orders',
version,
data: order,
})
for (const client of wss.clients) {
if (client.readyState === WebSocket.OPEN) {
client.send(message)
}
}
}, 5000)
console.log('WebSocket server listening on ws://localhost:8080')
运行:
npx tsx server.ts
前端开发环境连接地址应改为:
url: 'ws://localhost:8080'
预期行为:
- 浏览器建立 WebSocket 连接;
- 客户端发送
{ "type": "subscribe" }; - 服务端返回订单快照;
- 服务端每 5 秒广播一次
order.updated; - 客户端定期发送
ping; - 服务端返回
pong; - 停止服务端后,客户端进入关闭状态并按退避策略重连;
- 服务端恢复后,客户端重新订阅。
这个服务端示例没有实现事件历史,因此重连期间发生的更新可能丢失。要实现可靠恢复,需要增加事件日志、版本查询和快照降级逻辑。
十四、认证、代理与部署边界
14.1 认证
WebSocket 可以使用 Cookie,也可以在建立连接时携带协议允许的认证信息。浏览器原生 API 对自定义握手头的控制有限,不能假设可以像 fetch 那样任意设置 Authorization。
常见方案:
- 同源 Cookie + 服务端校验;
- 连接 URL 中使用短期一次性票据;
- 先通过 HTTPS 接口换取短期 WebSocket ticket;
- 由反向代理完成部分认证。
不要把长期访问令牌直接放在 WebSocket URL 中,因为 URL 可能出现在代理、监控和浏览器历史日志中。
SSE 的原生 EventSource 同样不提供任意自定义请求头。withCredentials: true 只解决跨源 Cookie 发送问题,不能替代 CORS、Cookie 的 SameSite、Secure 和服务端权限校验。
14.2 代理和负载均衡
WebSocket 需要代理正确转发升级请求;SSE 需要代理允许长时间保持响应,并及时刷新缓冲区。
故障表现可能是:
- 开发环境正常,生产环境立即
close; - SSE 连接看似建立,但数十秒没有事件;
- 消息集中到达,而不是实时到达;
- 连接被固定时间强制关闭。
诊断时应检查:
- 浏览器 Network 面板中的握手状态;
- WebSocket 是否返回
101 Switching Protocols; - SSE 是否返回
Content-Type: text/event-stream; - 代理是否启用了响应缓冲;
- 空闲超时是否短于心跳间隔;
- TLS、CORS、Cookie 和跨域策略;
- 服务端是否真正 flush 了事件。
14.3 多实例服务与广播
当服务端部署多个实例时,客户端可能连接到实例 A,业务写入发生在实例 B。若实例之间没有消息同步,客户端可能永远收不到更新。
通常需要:
- Redis Pub/Sub;
- 消息队列;
- Kafka 等事件流;
- 共享事件存储;
- 或让连接通过支持状态同步的网关转发。
负载均衡的会话粘性只能让同一个客户端尽量回到同一实例,不能自动解决实例之间的事件传播,也不能替代断线补偿。
十五、错误处理与诊断
15.1 将连接状态和同步状态分开
建议至少分别记录:
interface RealtimeDiagnostics {
connection: 'idle' | 'connecting' | 'open' | 'closed'
sync: 'unknown' | 'syncing' | 'live' | 'stale' | 'error'
reconnectAttempt: number
lastMessageAt: number | null
lastHeartbeatAt: number | null
lastPongAt: number | null
version: number
lastCloseCode: number | null
}
例如:
connection = open、sync = stale:连接还在,但本地数据有版本缺口;connection = closed、sync = stale:连接断开,等待恢复;connection = open、sync = live:连接存在且已完成同步。
仅展示“已连接”会掩盖真正的数据新鲜度问题。
15.2 记录关闭码和时间线
诊断日志至少应包含:
console.info('realtime event', {
event: 'close',
code: closeEvent.code,
reason: closeEvent.reason,
wasClean: closeEvent.wasClean,
attempt: reconnectAttempt.value,
lastMessageAt,
currentVersion,
})
同时记录:
- 连接开始时间;
- 握手成功时间;
- 最后一条消息时间;
- 最后一次 Pong 时间;
- 当前订阅;
- 当前版本;
- 重连等待时间;
- 是否主动关闭。
不要把完整业务数据和访问令牌直接写入日志。诊断字段应可关联,但要符合隐私和安全要求。
15.3 服务端错误与网络错误不同
服务端发送:
{
"type": "error",
"code": "AUTH_EXPIRED"
}
和浏览器触发 close 是两类错误:
- 业务错误需要客户端解释并决定是否刷新认证;
- 网络错误通常进入退避重连;
- 权限错误如果无限重连,只会制造无意义流量;
- 协议错误应停止当前连接并报警,而不是静默重试。
重连策略应根据错误类别决定,而不是所有错误都采用同一条路径。
十六、常见错误与反例
错误一:把 onerror 当作唯一重连入口
socket.onerror = () => {
socket = new WebSocket(url)
}
问题:
- 可能在连接尚未关闭时创建新连接;
- 旧连接仍有事件回调;
- 没有退避;
- 没有主动关闭判断;
- 可能产生多个并行连接。
应由 close 作为连接结束的统一入口,并确保一次只有一个重连定时器。
错误二:组件每次渲染都创建连接
const socket = new WebSocket(url)
如果这段代码位于响应式执行路径或未受生命周期控制,可能反复建立连接。连接创建应放入组合式函数的明确初始化逻辑中,并由 onBeforeUnmount 清理。
错误三:收到消息就直接覆盖数组
orders.value = message.data
如果 message.data 是增量事件,这会把单个对象误当成完整快照。必须根据 type 区分快照与增量:
if (message.type === 'snapshot') {
replaceAll(message.data)
} else if (message.type === 'order.updated') {
applyOne(message.data)
}
错误四:认为自动重连等于数据不丢失
SSE 的自动重连、WebSocket 的客户端重连都只解决连接恢复。没有事件历史、版本号或快照时,断线期间的数据仍可能丢失。
错误五:没有处理旧连接事件
关闭旧连接后立即创建新连接,如果旧连接回调仍能修改公共状态,就会出现状态跳回、重复消息和错误订阅。代次检查、事件监听器清理或连接控制器都可以解决这个问题。
错误六:只用“最后更新时间”判断新旧
客户端时间不是可靠的服务端顺序。应优先使用服务端版本、序列号或实体修订号;时间戳最多用于展示和诊断。
十七、测试实时系统的故障路径
不能只测试“服务端发送消息,页面显示消息”。至少需要验证:
17.1 连接路径
- 首次连接成功;
- URL、协议和认证失败;
- 服务端立即关闭;
- 连接一直处于
CONNECTING; - 页面销毁时连接被关闭。
17.2 重连路径
- 网络断开后是否进入退避;
- 是否只有一个重连定时器;
- 多次失败时延迟是否增长;
- 服务恢复后是否重置失败次数;
- 用户主动停止后是否不再重连。
17.3 心跳路径
- Ping 是否按预期发送;
- Pong 是否匹配对应请求;
- Pong 超时是否关闭连接;
- 页面后台时定时器节流是否导致误判;
- 服务端是否区分应用层 Ping 和业务消息。
17.4 同步路径
- 重复版本是否被忽略;
- 版本缺口是否触发补偿;
- 快照是否替换旧状态;
- 重连后是否带上正确的
lastVersion; - 事件历史过期时是否能回退到快照;
- 本地乐观更新与服务端事件冲突时结果是否明确。
浏览器开发者工具可以模拟 Offline,但它不一定覆盖所有真实的半开连接、代理超时和移动网络切换场景。服务端也应提供测试开关,模拟延迟、乱序、重复、断开和错误消息。
十八、生产取舍:可靠性等级要明确
实时数据常见的投递语义包括:
至多一次
消息发送一次,失败后不保证重试。
优点是简单、延迟低;缺点是容易丢消息。适合不重要的在线状态或瞬时提示。
至少一次
消息可能被重复发送,但尽量不丢失。
客户端和服务端必须支持幂等、去重和版本检查。许多断线重连系统实际采用这种模型。
恰好一次
每条业务操作严格只生效一次。
这通常不是单靠 WebSocket 就能实现的,需要:
- 服务端幂等键;
- 持久化命令状态;
- 明确的事务边界;
- 客户端确认和重试协议;
- 去重存储。
因此,不应把“连接稳定”误认为“业务操作恰好执行一次”。
对于大多数 Vue 实时页面,一个合理的基础模型是:
服务端权威状态
+ 版本号
+ 快照
+ 增量事件
+ 至少一次投递
+ 客户端去重
+ 版本缺口时重新同步
这个模型比单纯依赖连接状态更可靠,也比试图在前端自行推导全部状态更容易验证。
十九、最后的实现边界
WebSocket 解决双向长连接,SSE 解决服务端事件流;重连解决连接暂时失效,心跳解决空闲或半开连接探测,状态同步解决断线、重复、乱序和版本缺口。它们是不同层次的问题,不能互相替代。
在 Vue 应用中,较完整的数据流应是:
组件生命周期或 Pinia
↓
连接管理
↓
重连与心跳
↓
消息校验
↓
版本检查、去重、补偿
↓
权威状态写入 Pinia
↓
组件响应式渲染
如果只实现了 new WebSocket() 和 onmessage,得到的是“能接收消息的页面”;只有当连接生命周期、故障恢复和状态版本都被定义后,才形成可验证的实时数据系统。
系列导航与关联阅读
- 系列入口:Vue 完整学习路线:从响应式与组件到工程化、SSR 和生产交付
- 上一篇:Vue HTTP 客户端封装:Axios、拦截器、取消、重试和错误模型
- 下一篇:Vue 单元测试:Vitest、Composable、时间、网络和稳定断言
官方资料
本文依据 Vue、Vite 与生态项目官方文档重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论