Vue 基础体系 · 第 48/70 篇。示例基于 Vue 3、Composition API、TypeScript 与现代 Vite 工具链;版本敏感能力会单独标注。

Vue 实时数据:WebSocket、SSE、重连、心跳和状态同步

实时数据不是“建立一条长连接,然后在 onmessage 中改几个变量”这么简单。真正的实时系统至少要回答五个问题:

  1. 客户端如何与服务端建立双向或单向连接?
  2. 连接断开后,什么时候重连,重连多少次,如何避免连接风暴?
  3. 如何判断连接只是暂时空闲,还是已经不可用?
  4. 重连期间丢失的数据如何恢复?
  5. 多条消息并发到达、重复到达或乱序到达时,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 openmessageerrorclose

浏览器 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 的主要特点是:

  • 浏览器原生支持自动重连;
  • 服务端可以发送命名事件;
  • 客户端通过 messageaddEventListener 接收消息;
  • 客户端不能通过同一个 EventSource 连接发送业务消息;
  • 若客户端需要发送命令,仍需使用 fetch、普通 HTTP 请求或另一个 WebSocket。

3.2 SSE 的 idLast-Event-ID

SSE 事件可以携带 id

id: 42
data: {"status":"paid"}

浏览器会记录最近收到的事件 ID。连接意外断开后,浏览器通常会在重连请求中发送:

Last-Event-ID: 42

服务端可以根据这个 ID 补发后续事件。

这提供了“从上次位置继续”的基础,但它不是完整的数据一致性保证。服务端仍然必须:

  1. 保存一段时间的事件历史;
  2. 判断客户端的 ID 是否仍在可补发范围内;
  3. 超出范围时发送完整快照;
  4. 保证事件 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 如何选择

可以用一个直接的判断过程:

  1. 如果只有服务端向客户端推送,先考虑 SSE。
  2. 如果客户端也需要低延迟发送消息,考虑 WebSocket。
  3. 如果需要浏览器原生自动重连,并且消息是事件流,SSE 更简单。
  4. 如果需要二进制帧、双向请求响应、订阅和取消订阅,WebSocket 更合适。
  5. 如果部署环境对长 HTTP 响应支持不稳定,SSE 需要重点验证代理和网关配置。
  6. 如果客户端需要带复杂认证头,原生 EventSource 可能不合适。

二者都不能自动解决状态同步问题。协议只负责传输,数据是否完整、是否重复、是否乱序,需要应用协议负责


五、重连:从“重新连接”到可控算法

5.1 为什么不能立即无限重连

假设服务端宕机,客户端在 close 后立即重连:

function reconnect() {
  socket = new WebSocket(url)
}

如果每个客户端都这样做,结果会是:

  1. 所有连接几乎同时失败;
  2. 所有客户端几乎同时发起新连接;
  3. 服务端仍未恢复,连接再次失败;
  4. 客户端进入高频重试;
  5. 服务端恢复时遭遇连接洪峰。

这种现象称为连接风暴。解决它需要退避和抖动。

5.2 指数退避与抖动

设第 nn 次重连等待时间为:

dn=min(dmax,d02n)+Jnd_n = \min(d_{\max}, d_0 \cdot 2^n) + J_n

其中:

  • d0d_0 是初始延迟,例如 1000 毫秒;
  • dmaxd_{\max} 是最大延迟,例如 30 秒;
  • nn 是已经失败的重连次数;
  • JnJ_n 是随机抖动;
  • 常见全抖动方式为 Jn[0,dn]J_n \in [0, d_n] 的随机值。

完整算法可以写成:

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 = 1000max = 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

假设发生以下顺序:

  1. 连接 A 断开;
  2. 客户端创建连接 B;
  3. 连接 A 的某个延迟事件才到达;
  4. 连接 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,无法发现半开连接。完整逻辑应包括:

  1. 每隔 TT 毫秒发送 Ping;
  2. 记录本次 Ping 的时间和请求 ID;
  3. HH 毫秒内等待 Pong;
  4. 超时后主动关闭连接;
  5. 由关闭流程触发重连。

可以定义:

连接健康=收到对应 Pong(当前时间发送时间H)\text{连接健康} = \text{收到对应 Pong} \land (\text{当前时间} - \text{发送时间} \leq H)

例如:

  • 心跳间隔 T=20T = 20 秒;
  • Pong 超时 H=10H = 10 秒。

若客户端连续发送 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
}

这里使用了三个条件:

  1. version <= currentVersion:旧消息或重复消息;
  2. version === currentVersion + 1:正好是下一条;
  3. 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 会造成中间状态不确定。

有三种处理策略:

  1. 严格连续策略:发现版本不是 current + 1 就暂停并重新同步。
  2. 客户端缓冲策略:暂存高版本消息,等待缺失版本到达。
  3. 按实体版本策略:每个实体独立比较版本,只应用该实体更新版本更高的消息。

选择取决于业务:

  • 金融余额、库存、权限等状态通常需要严格连续或服务端快照;
  • 独立商品卡片可以按实体版本更新;
  • 聊天消息可能需要全局序列保证显示顺序。

12.3 本地操作与服务端推送并发

假设用户点击“支付”:

  1. 客户端立即把订单显示为 paid
  2. 服务端稍后推送版本 50,状态仍是 pending
  3. 客户端如果无条件覆盖,就会出现状态回退。

常见解决方案有:

  • 不做乐观更新,等待服务端确认;
  • 为本地命令分配 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'

预期行为:

  1. 浏览器建立 WebSocket 连接;
  2. 客户端发送 { "type": "subscribe" }
  3. 服务端返回订单快照;
  4. 服务端每 5 秒广播一次 order.updated
  5. 客户端定期发送 ping
  6. 服务端返回 pong
  7. 停止服务端后,客户端进入关闭状态并按退避策略重连;
  8. 服务端恢复后,客户端重新订阅。

这个服务端示例没有实现事件历史,因此重连期间发生的更新可能丢失。要实现可靠恢复,需要增加事件日志、版本查询和快照降级逻辑。


十四、认证、代理与部署边界

14.1 认证

WebSocket 可以使用 Cookie,也可以在建立连接时携带协议允许的认证信息。浏览器原生 API 对自定义握手头的控制有限,不能假设可以像 fetch 那样任意设置 Authorization

常见方案:

  • 同源 Cookie + 服务端校验;
  • 连接 URL 中使用短期一次性票据;
  • 先通过 HTTPS 接口换取短期 WebSocket ticket;
  • 由反向代理完成部分认证。

不要把长期访问令牌直接放在 WebSocket URL 中,因为 URL 可能出现在代理、监控和浏览器历史日志中。

SSE 的原生 EventSource 同样不提供任意自定义请求头。withCredentials: true 只解决跨源 Cookie 发送问题,不能替代 CORS、Cookie 的 SameSiteSecure 和服务端权限校验。

14.2 代理和负载均衡

WebSocket 需要代理正确转发升级请求;SSE 需要代理允许长时间保持响应,并及时刷新缓冲区。

故障表现可能是:

  • 开发环境正常,生产环境立即 close
  • SSE 连接看似建立,但数十秒没有事件;
  • 消息集中到达,而不是实时到达;
  • 连接被固定时间强制关闭。

诊断时应检查:

  1. 浏览器 Network 面板中的握手状态;
  2. WebSocket 是否返回 101 Switching Protocols
  3. SSE 是否返回 Content-Type: text/event-stream
  4. 代理是否启用了响应缓冲;
  5. 空闲超时是否短于心跳间隔;
  6. TLS、CORS、Cookie 和跨域策略;
  7. 服务端是否真正 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 = opensync = stale:连接还在,但本地数据有版本缺口;
  • connection = closedsync = stale:连接断开,等待恢复;
  • connection = opensync = 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、Vite 与生态项目官方文档重新梳理;正文与示例由 WR BLOG 编写。