CC 咖啡猫的工作空间 Coding Space

实时通信实践(SSE + WebSocket)

在现代 Web 应用中,实时通信已从"加分项"变成了"基础能力"。无论是即时消息、协作编辑、数据大屏,还是 AI 流式输出,都离不开实时通信技术。

本文将系统性地讨论四种主流的实时通信方案(短轮询、长轮询、SSE、WebSocket),并重点展开 WebSocket 和 SSE 的最佳实践。


1. 实时通信方案对比

1.1 短轮询(Short Polling)

原理:前端通过 setIntervalsetTimeout 定时向服务端发起 HTTP 请求,询问是否有新数据。

Client                          Server
  |                               |
  |--- GET /api/messages?since=1  |
  |← 200 OK { data: [...] }      |
  |--- GET /api/messages?since=2  |
  |← 200 OK { data: [] }         |
  |--- GET /api/messages?since=2  |
  |← 200 OK { data: [3] }        |
  |                               |

适用场景:对实时性要求极低(如 30s 以上刷新间隔)的后台管理页面。

维度 评价
实时性 差,取决于轮询间隔(通常 3s~30s)
服务端压力 高,大量无意义的 HTTP 请求
带宽消耗 高,每次请求都带完整的 HTTP 头部
实现复杂度 极低
浏览器兼容性 最好,纯 HTTP
// ❌ 短轮询示例 — 不推荐用于实时场景

function startPolling() {
  setInterval(async () => {
    try {
      const res = await fetch('/api/messages');
      const data = await res.json();
      // 即使没有新消息,也发起了一次完整的 HTTP 请求
      if (data?.messages?.length) {
        // 处理新消息
      }
    } catch (error) {
      console.error('Polling failed:', error);
    }
  }, 3000); // 3 秒一次
}

1.2 长轮询(Long Polling)

原理:客户端发起 HTTP 请求到服务端,服务端在没有新数据时挂起请求(不立即返回),直到有新数据时才返回响应。客户端收到响应后,立即发起下一次请求。

Client                          Server
  |                               |
  |--- GET /api/poll?since=1      |
  |       (请求被挂起,等待新消息)    |
  |              ...              |
  |← 200 OK { data: [2, 3] }     | ← 有新消息,立即返回
  |--- GET /api/poll?since=3      |
  |       (请求被挂起,等待新消息)    |
  |                               |

适用场景:实时性要求中等、无法使用 WebSocket 的旧版浏览器环境。

维度 评价
实时性 较好(消息到达即返回)
服务端压力 中等(挂起连接数受限于服务器并发)
带宽消耗 中等(每次仍带完整 HTTP 头部)
实现复杂度 中等(需要服务端配合挂起逻辑)
兼容性 好(纯 HTTP)
// ❌ 长轮询示例 — 有更好的替代方案

async function longPolling(since: number) {
  try {
    const res = await fetch(`/api/poll?since=${since}`);
    const data = await res.json();
    if (data?.messages?.length) {
      // 处理消息
      since = data.messages[data.messages.length - 1].id;
    }
  } catch {
    // 重连逻辑
  } finally {
    // 无论成功还是失败,立即发起下一次长轮询
    longPolling(since);
  }
}

1.3 SSE(Server-Sent Events)

原理:通过 HTTP 协议建立一条持久连接,服务端可以持续向客户端推送数据。使用 EventSource API 原生支持。

Client                          Server
  |                               |
  |--- GET /api/events            |
  |← 200 OK Content-Type: text/event-stream
  |← data: {"type": "message", "content": "hello"}
  |← data: {"type": "notification", "count": 5}
  |← data: {"type": "ping"}
  |                               |

适用场景:服务端单向推送(如通知、数据大屏、AI 流式输出)。

1.4 WebSocket

原理:基于 TCP 协议的全双工通信协议,通过 HTTP 升级握手建立连接,之后双方可以自由发送数据。

Client                          Server
  |                               |
  |--- Upgrade: websocket         |
  |← 101 Switching Protocols     |
  |=== 全双工通信 ================|
  |--- { type: "message", content: "hello" }
  |← { type: "ack", id: 1 }
  |← { type: "notification" }
  |--- { type: "ping" }
  |                               |

适用场景:双向实时通信(IM 聊天、协作编辑、游戏)。

1.5 方案对比总结

维度 短轮询 长轮询 SSE WebSocket
通信方向 客户端请求 客户端请求 服务端推送 双向
实时性 低(取决于间隔) 最高
协议 HTTP/1.0+ HTTP/1.0+ HTTP + text/event-stream WS / WSS
服务端压力 高(大量空请求) 中(连接挂起) 低(一条连接) 低(一条连接)
自动重连 内置 无(需自行实现)
二进制支持 否(纯文本)
跨域支持 CORS CORS CORS(有限制) 无限制
浏览器兼容性 全部 全部 IE 不支持 IE10+
实现复杂度 极低
适用场景 非实时场景 兼容旧浏览器 通知/流式输出/大屏 即时通讯/游戏/协作

选型建议

  • 服务端单向推送 => SSE(实现简单,自带重连)
  • 双向实时通信 => WebSocket
  • 需要兼容老旧浏览器 => 长轮询(SSE/WebSocket 兜底降级)
  • 不需要实时,只是定时刷新 => 短轮询

2. WebSocket 实践

2.1 原生 WebSocket API 封装

原生 WebSocket API 功能较弱,需要手动封装心跳、重连、ACK 等能力。

// src/utils/websocket.ts — 原生 WebSocket 封装

import { EventEmitter } from 'eventemitter3'

type ConnectionState = 'connecting' | 'open' | 'closing' | 'closed'

interface WebSocketOptions {
  url: string
  protocols?: string | string[]
  reconnectInterval?: number     // 重连间隔(ms),默认 1000
  maxReconnectAttempts?: number  // 最大重连次数,默认 10
  heartbeatInterval?: number     // 心跳间隔(ms),默认 30000
  heartbeatMessage?: object      // 心跳消息内容
  onOpen?: () => void
  onClose?: (event: CloseEvent) => void
  onError?: (error: Event) => void
  onMessage?: (data: unknown) => void
}

class WebSocketClient extends EventEmitter {
  private ws: WebSocket | null = null
  private options: Required<WebSocketOptions>
  private state: ConnectionState = 'closed'
  private reconnectAttempts = 0
  private heartbeatTimer: ReturnType<typeof setInterval> | null = null
  private reconnectTimer: ReturnType<typeof setTimeout> | null = null
  private manuallyClosed = false

  constructor(options: WebSocketOptions) {
    super()

    // 防御式编程:提供完整的默认值
    this.options = {
      url: options.url,
      protocols: options.protocols,
      reconnectInterval: options.reconnectInterval ?? 1000,
      maxReconnectAttempts: options.maxReconnectAttempts ?? 10,
      heartbeatInterval: options.heartbeatInterval ?? 30000,
      heartbeatMessage: options.heartbeatMessage ?? { type: 'ping' },
      onOpen: options.onOpen ?? (() => {}),
      onClose: options.onClose ?? (() => {}),
      onError: options.onError ?? (() => {}),
      onMessage: options.onMessage ?? (() => {}),
    }

    this.connect()
  }

  /** 获取当前连接状态 */
  get connectionState(): ConnectionState {
    return this.state
  }

  /** 连接 WebSocket */
  private connect() {
    if (this.ws?.readyState === WebSocket.OPEN) return

    this.state = 'connecting'
    this.emit('stateChange', this.state)

    try {
      this.ws = this.options.protocols
        ? new WebSocket(this.options.url, this.options.protocols)
        : new WebSocket(this.options.url)
    } catch (error) {
      console.error('[WebSocket] 创建连接失败:', error)
      this.handleReconnect()
      return
    }

    this.ws.onopen = () => {
      this.state = 'open'
      this.reconnectAttempts = 0
      this.emit('stateChange', this.state)
      this.options.onOpen()
      this.startHeartbeat()
    }

    this.ws.onclose = (event: CloseEvent) => {
      this.state = 'closed'
      this.emit('stateChange', this.state)
      this.options.onClose(event)
      this.stopHeartbeat()

      if (!this.manuallyClosed) {
        this.handleReconnect()
      }
    }

    this.ws.onerror = (error: Event) => {
      console.error('[WebSocket] 连接错误:', error)
      this.options.onError(error)
    }

    this.ws.onmessage = (event: MessageEvent) => {
      try {
        // 防御式编程:始终校验消息格式
        const data = JSON.parse(event.data) as Record<string, unknown>
        this.options.onMessage(data)
        this.emit('message', data)
      } catch {
        // 非 JSON 格式的消息,直接透传
        this.options.onMessage(event.data)
        this.emit('message', event.data)
      }
    }
  }

  /** 发送消息 */
  send(data: unknown): boolean {
    if (this.ws?.readyState !== WebSocket.OPEN) {
      console.warn('[WebSocket] 连接未就绪,无法发送消息')
      return false
    }

    try {
      this.ws.send(typeof data === 'string' ? data : JSON.stringify(data))
      return true
    } catch (error) {
      console.error('[WebSocket] 发送消息失败:', error)
      return false
    }
  }

  /** 关闭连接 */
  close() {
    this.manuallyClosed = true
    this.stopHeartbeat()
    this.clearReconnectTimer()
    this.ws?.close()
  }

  /** 重新连接(指数退避) */
  private handleReconnect() {
    if (this.reconnectAttempts >= this.options.maxReconnectAttempts) {
      console.error('[WebSocket] 已达最大重连次数,停止重连')
      this.emit('maxReconnectExceeded')
      return
    }

    // 指数退避:1s → 2s → 4s → 8s → ...(最大 30s)
    const delay = Math.min(
      this.options.reconnectInterval * Math.pow(2, this.reconnectAttempts),
      30000,
    )

    console.log(`[WebSocket] ${delay}ms 后尝试第 ${this.reconnectAttempts + 1} 次重连`)
    this.reconnectTimer = setTimeout(() => {
      this.reconnectAttempts++
      this.connect()
    }, delay)
  }

  /** 启动心跳 */
  private startHeartbeat() {
    this.stopHeartbeat()
    this.heartbeatTimer = setInterval(() => {
      this.send(this.options.heartbeatMessage)
    }, this.options.heartbeatInterval)
  }

  /** 停止心跳 */
  private stopHeartbeat() {
    if (this.heartbeatTimer) {
      clearInterval(this.heartbeatTimer)
      this.heartbeatTimer = null
    }
  }

  /** 清除重连定时器 */
  private clearReconnectTimer() {
    if (this.reconnectTimer) {
      clearTimeout(this.reconnectTimer)
      this.reconnectTimer = null
    }
  }
}

export { WebSocketClient, type WebSocketOptions, type ConnectionState }

2.2 Socket.IO 推荐

Socket.IO 是目前最成熟的 WebSocket 封装库,提供了原生 WebSocket 所缺乏的诸多能力。

特性 原生 WebSocket Socket.IO
自动重连 需要自行实现 内置,支持退避策略
心跳/Ping 需要自行实现 内置
事件分发 统一 onmessage 按事件名分发 (socket.on('event', cb))
房间 原生支持
ACK 确认 需要自行实现 内置回调机制
降级 不支持 不支持 WebSocket 时自动降级到 HTTP 长轮询
二进制 支持 支持
包体积 0(浏览器内置) ~50KB(gzip)
服务端 任意语言 需要 Node.js 服务端
# 安装
pnpm add socket.io-client

Socket.IO 基础用法

// src/utils/socket.ts — Socket.IO 客户端封装

import { io, Socket } from 'socket.io-client'

let socket: Socket | null = null

interface SocketOptions {
  url: string
  token?: string
  path?: string
  reconnectionAttempts?: number
}

export function createSocket(options: SocketOptions): Socket {
  if (socket?.connected) {
    return socket
  }

  socket = io(options.url, {
    path: options.path ?? '/socket.io',
    auth: options.token ? { token: options.token } : undefined,
    transports: ['websocket', 'polling'], // 首选 WebSocket,降级到 polling
    reconnection: true,
    reconnectionAttempts: options.reconnectionAttempts ?? 10,
    reconnectionDelay: 1000,
    reconnectionDelayMax: 30000,
    randomizationFactor: 0.5,
    timeout: 10000,
  })

  socket.on('connect', () => {
    console.log('[Socket.IO] 已连接, id:', socket?.id)
  })

  socket.on('disconnect', (reason) => {
    console.log('[Socket.IO] 断开连接:', reason)
    // 'io server disconnect' => 服务端主动断开,不应重连
    // 'transport close' / 'transport error' => 网络问题,会触发自动重连
  })

  socket.on('connect_error', (error) => {
    console.error('[Socket.IO] 连接错误:', error.message)
  })

  return socket
}

export function getSocket(): Socket | null {
  return socket
}

export function disconnectSocket() {
  socket?.disconnect()
  socket = null
}

2.3 Vue:useWebSocket Composable

// src/composables/useWebSocket.ts — Vue 3 Composable

import { ref, computed, onUnmounted, type Ref } from 'vue'
import { WebSocketClient, type WebSocketOptions } from '@/utils/websocket'

export function useWebSocket(options: WebSocketOptions) {
  const connectionState = ref<'connecting' | 'open' | 'closing' | 'closed'>('closed')
  const lastMessage = ref<unknown>(null)
  const reconnectCount = ref(0)

  let client: WebSocketClient | null = null

  /** 连接状态描述 */
  const statusText = computed(() => {
    const map: Record<string, string> = {
      connecting: '连接中...',
      open: '已连接',
      closing: '关闭中...',
      closed: '未连接',
    }
    return map[connectionState.value] || '未知'
  })

  /** 是否已连接 */
  const isConnected = computed(() => connectionState.value === 'open')

  function connect() {
    client = new WebSocketClient({
      ...options,
      onOpen() {
        connectionState.value = 'open'
        options.onOpen?.()
      },
      onClose(event) {
        connectionState.value = 'closed'
        options.onClose?.(event)
      },
      onError(error) {
        options.onError?.(error)
      },
      onMessage(data) {
        lastMessage.value = data
        options.onMessage?.(data)
      },
    })

    client.on('stateChange', (state: string) => {
      connectionState.value = state as 'connecting' | 'open' | 'closing' | 'closed'
    })

    client.on('maxReconnectExceeded', () => {
      reconnectCount.value++
    })
  }

  function send(data: unknown): boolean {
    return client?.send(data) ?? false
  }

  function disconnect() {
    client?.close()
    client = null
  }

  onUnmounted(() => {
    disconnect()
  })

  return {
    connectionState,
    statusText,
    isConnected,
    lastMessage,
    reconnectCount,
    connect,
    send,
    disconnect,
  }
}
<!-- src/components/ChatRoom.vue — Vue 3 + Element Plus + WebSocket -->
<script setup lang="ts">
import { ref, onMounted } from 'vue'
import { ElInput, ElButton, ElAlert, ElBadge } from 'element-plus'
import { useWebSocket } from '@/composables/useWebSocket'

const { connectionState, statusText, isConnected, send } = useWebSocket({
  url: import.meta.env.VITE_WS_URL,
  heartbeatMessage: { type: 'ping' },
  onMessage(data) {
    const msg = data as { type: string; content: string; sender: string }
    if (msg?.type === 'message') {
      messages.value.push(msg)
    }
  },
})

const messages = ref<Array<{ type: string; content: string; sender: string }>>([])
const inputText = ref('')

function handleSend() {
  const text = inputText.value.trim()
  if (!text) return

  // 防御式编程:连接未就绪时禁止发送
  if (!isConnected.value) {
    ElMessage.warning('连接未就绪,请等待重连...')
    return
  }

  const sent = send({
    type: 'message',
    content: text,
    timestamp: Date.now(),
  })

  if (!sent) {
    ElMessage.error('发送失败')
    return
  }

  inputText.value = ''
}

onMounted(() => {
  // 自动连接
})
</script>

<template>
  <div class="chat-room">
    <ElAlert
      :title="statusText"
      :type="isConnected ? 'success' : 'warning'"
      :closable="false"
      show-icon
    />
    <div class="message-list">
      <div v-for="(msg, index) in messages" :key="index">
        <strong>{{ msg.sender }}: </strong>{{ msg.content }}
      </div>
    </div>
    <div class="input-area">
      <ElInput
        v-model="inputText"
        placeholder="输入消息..."
        @keyup.enter="handleSend"
      />
      <ElButton type="primary" @click="handleSend" :disabled="!isConnected">
        发送
      </ElButton>
    </div>
  </div>
</template>

2.4 React:useWebSocket Hook

// src/hooks/useWebSocket.ts — React Hook

import { useState, useEffect, useRef, useCallback } from 'react'
import { WebSocketClient, type WebSocketOptions } from '@/utils/websocket'

type ConnectionState = 'connecting' | 'open' | 'closing' | 'closed'

interface UseWebSocketReturn {
  connectionState: ConnectionState
  isConnected: boolean
  statusText: string
  lastMessage: unknown
  reconnectCount: number
  send: (data: unknown) => boolean
  disconnect: () => void
}

export function useWebSocket(options: WebSocketOptions): UseWebSocketReturn {
  const [connectionState, setConnectionState] = useState<ConnectionState>('closed')
  const [lastMessage, setLastMessage] = useState<unknown>(null)
  const [reconnectCount, setReconnectCount] = useState(0)
  const clientRef = useRef<WebSocketClient | null>(null)

  const statusText = {
    connecting: '连接中...',
    open: '已连接',
    closing: '关闭中...',
    closed: '未连接',
  }[connectionState]

  const isConnected = connectionState === 'open'

  useEffect(() => {
    const client = new WebSocketClient({
      ...options,
      onOpen() {
        setConnectionState('open')
        options.onOpen?.()
      },
      onClose(event) {
        setConnectionState('closed')
        options.onClose?.(event)
      },
      onError(error) {
        options.onError?.(error)
      },
      onMessage(data) {
        setLastMessage(data)
        options.onMessage?.(data)
      },
    })

    client.on('stateChange', (state: string) => {
      setConnectionState(state as ConnectionState)
    })

    client.on('maxReconnectExceeded', () => {
      setReconnectCount((prev) => prev + 1)
    })

    clientRef.current = client

    return () => {
      client.close()
    }
  }, [options.url])

  const send = useCallback((data: unknown): boolean => {
    return clientRef.current?.send(data) ?? false
  }, [])

  const disconnect = useCallback(() => {
    clientRef.current?.close()
    clientRef.current = null
  }, [])

  return {
    connectionState,
    isConnected,
    statusText,
    lastMessage,
    reconnectCount,
    send,
    disconnect,
  }
}
// src/components/ChatRoom.tsx — React + Ant Design + WebSocket

import { useState } from 'react'
import { Input, Button, Alert, message } from 'antd'
import { useWebSocket } from '@/hooks/useWebSocket'

interface ChatMessage {
  type: string
  content: string
  sender: string
  timestamp?: number
}

const ChatRoom: React.FC = () => {
  const { connectionState, isConnected, statusText, send } = useWebSocket({
    url: import.meta.env.VITE_WS_URL,
    heartbeatMessage: { type: 'ping' },
    onMessage(data) {
      const msg = data as ChatMessage
      if (msg?.type === 'message') {
        setMessages((prev) => [...prev, msg])
      }
    },
  })

  const [messages, setMessages] = useState<ChatMessage[]>([])
  const [inputText, setInputText] = useState('')

  const handleSend = () => {
    const text = inputText.trim()
    if (!text) return

    // 防御式编程:连接未就绪时禁止发送
    if (!isConnected) {
      message.warning('连接未就绪,请等待重连...')
      return
    }

    const sent = send({
      type: 'message',
      content: text,
      timestamp: Date.now(),
    })

    if (!sent) {
      message.error('发送失败')
      return
    }

    setInputText('')
  }

  return (
    <div className="chat-room">
      <Alert
        message={statusText}
        type={isConnected ? 'success' : 'warning'}
        showIcon
        closable={false}
      />
      <div className="message-list">
        {messages.map((msg, index) => (
          <div key={index}>
            <strong>{msg.sender}: </strong>
            {msg.content}
          </div>
        ))}
      </div>
      <div className="input-area">
        <Input.TextArea
          value={inputText}
          onChange={(e) => setInputText(e.target.value)}
          onPressEnter={handleSend}
          placeholder="输入消息..."
          rows={2}
        />
        <Button
          type="primary"
          onClick={handleSend}
          disabled={!isConnected}
        >
          发送
        </Button>
      </div>
    </div>
  )
}

export default ChatRoom

2.5 心跳检测(Ping / Pong)

WebSocket 连接可能因为网络问题、NAT 超时、代理断开等原因悄然断开,而 onclose 事件可能不会被触发。心跳检测是保活的必要手段。

机制

  • 客户端定时发送 Ping(通常 30s)
  • 服务端收到 Ping 后回复 Pong
  • 如果客户端超过 N 秒未收到 Pong,认为连接断开,主动关闭并触发重连
间隔配置 推荐值 说明
心跳发送间隔 30s 大多数云环境 NAT 超时是 5 分钟
心跳超时 10s 发送后未收到 Pong 的等待时间
重连间隔(最小) 1s 初始重连等待时间
重连间隔(最大) 30s 指数退避上限
// Socket.IO 内置心跳,无需手动实现
const socket = io(url, {
  pingInterval: 25000,  // 25s 发送一次 ping
  pingTimeout: 20000,   // 20s 未收到 pong 视为断开
})
// 原生 WebSocket 心跳封装(上述 WebSocketClient 已内置)
// 关键代码片段:

private startHeartbeat() {
  this.stopHeartbeat()
  this.heartbeatTimer = setInterval(() => {
    this.send(this.options.heartbeatMessage)
    // 防御式编程:如果连接已断开但定时器未清理
    if (this.ws?.readyState !== WebSocket.OPEN) {
      this.stopHeartbeat()
    }
  }, this.options.heartbeatInterval)
}

2.6 断线重连(指数退避)

指数退避(Exponential Backoff)是断线重连的标准策略,避免客户端在服务端故障恢复期间集中重连造成雪崩。

算法

delay = min(baseDelay * 2^attempt, maxDelay)
重连次数 延迟时间
0 1s
1 2s
2 4s
3 8s
4 16s
5+ 30s(上限)
// 指数退避重连实现

const BASE_DELAY = 1000    // 1s
const MAX_DELAY = 30000    // 30s
const MAX_RETRIES = 10

function calculateDelay(attempt: number): number {
  // 指数退避 + 随机抖动,防止同时重连的 TCP 同步风暴
  const exponentialDelay = Math.min(BASE_DELAY * Math.pow(2, attempt), MAX_DELAY)
  // 随机抖动 ±50%
  const jitter = exponentialDelay * (0.5 + Math.random() * 0.5)
  return Math.floor(jitter)
}

// 使用
for (let attempt = 0; attempt < MAX_RETRIES; attempt++) {
  const delay = calculateDelay(attempt)
  console.log(`[WebSocket] 将在 ${delay}ms 后尝试第 ${attempt + 1} 次重连`)
  await new Promise((resolve) => setTimeout(resolve, delay))
  // ... 尝试连接
}

2.7 消息确认(ACK)

在实时通信中,消息可能因为网络问题丢失。ACK(Acknowledgment)机制确保消息被服务端正确接收。

流程

Client                          Server
  |                               |
  |--- { id: "msg_001", type: "message", content: "hello" }
  |                               | 处理消息
  |← { id: "msg_001", type: "ack" }
  |                               |
  |--- { id: "msg_002", type: "message", content: "world" }
  |← { id: "msg_002", type: "ack" }
  |                               |
// src/utils/ack-manager.ts — ACK 管理器

interface PendingMessage {
  id: string
  data: unknown
  timestamp: number
  retries: number
  timer: ReturnType<typeof setTimeout>
}

class AckManager {
  private pending = new Map<string, PendingMessage>()
  private maxRetries = 3
  private timeout = 10000 // 10s 未收到 ACK 则重发
  private onResend: (data: unknown) => boolean

  constructor(onResend: (data: unknown) => boolean) {
    this.onResend = onResend
  }

  /** 生成唯一消息 ID */
  generateId(): string {
    return `msg_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`
  }

  /** 发送消息并等待 ACK */
  send(id: string, data: unknown, sendFn: () => boolean): boolean {
    if (!sendFn()) return false

    // 添加到待确认队列
    const timer = setTimeout(() => this.handleTimeout(id), this.timeout)
    this.pending.set(id, { id, data, timestamp: Date.now(), retries: 0, timer })
    return true
  }

  /** 处理 ACK 响应 */
  handleAck(id: string) {
    const pending = this.pending.get(id)
    if (pending) {
      clearTimeout(pending.timer)
      this.pending.delete(id)
    }
  }

  /** 超时重试 */
  private handleTimeout(id: string) {
    const pending = this.pending.get(id)
    if (!pending) return

    if (pending.retries >= this.maxRetries) {
      console.error(`[ACK] 消息 ${id} 已达最大重试次数,放弃发送`)
      this.pending.delete(id)
      return
    }

    pending.retries++
    const resendSuccess = this.onResend(pending.data)

    if (resendSuccess) {
      // 重新计时
      const timer = setTimeout(() => this.handleTimeout(id), this.timeout)
      pending.timer = timer
    }
  }

  /** 清理所有待确认消息 */
  clear() {
    this.pending.forEach((pending) => clearTimeout(pending.timer))
    this.pending.clear()
  }
}

export { AckManager }
// 配合 Socket.IO 使用 ACK(Socket.IO 内置支持)

// 客户端发送并等待 ACK
socket.emit('message', { content: 'hello' }, (ack: { id: string; status: string }) => {
  // 这个回调在服务端处理完成后触发
  console.log('服务端确认:', ack)
})

// 服务端回复 ACK
socket.on('message', (data, callback) => {
  // 处理消息...
  callback({ id: data.id, status: 'received' }) // 相当于发送 ACK
})

2.8 连接状态管理

连接状态管理是实时应用的基石,状态变化直接影响 UI 展示和用户交互。

状态 含义 用户提示 UI 表现
connecting 正在建立连接 "连接中..." 加载中状态,发送按钮禁用
open 连接已建立 "已连接" 正常聊天,绿色状态灯
closing 正在关闭 "关闭中..." 短暂过渡状态
closed 连接已断开 "未连接" 红色状态灯,显示重连提示
// 连接状态在状态管理中的统一管理

// Vue + Pinia
// src/stores/websocket.ts
import { defineStore } from 'pinia'
import { ref } from 'vue'

export const useWebSocketStore = defineStore('websocket', () => {
  const connectionState = ref<'connecting' | 'open' | 'closing' | 'closed'>('closed')
  const reconnectAttempts = ref(0)
  const serverTimeDiff = ref(0) // 服务端与客户端时间差,用于消息排序

  function updateState(state: 'connecting' | 'open' | 'closing' | 'closed') {
    connectionState.value = state
    if (state === 'open') {
      reconnectAttempts.value = 0
    }
  }

  return { connectionState, reconnectAttempts, serverTimeDiff, updateState }
})
// React + Zustand
// src/stores/websocket.ts
import { create } from 'zustand'

type ConnectionState = 'connecting' | 'open' | 'closing' | 'closed'

interface WebSocketStore {
  connectionState: ConnectionState
  reconnectAttempts: number
  serverTimeDiff: number
  updateState: (state: ConnectionState) => void
}

export const useWebSocketStore = create<WebSocketStore>((set) => ({
  connectionState: 'closed',
  reconnectAttempts: 0,
  serverTimeDiff: 0,
  updateState: (state) =>
    set({
      connectionState: state,
      reconnectAttempts: state === 'open' ? 0 : undefined,
    }),
}))

3. SSE 实践

3.1 EventSource API 使用

SSE(Server-Sent Events)的核心是浏览器原生 EventSource API,它比 WebSocket 简单得多。

// 基础用法
const eventSource = new EventSource('/api/stream')

// 监听命名事件
eventSource.addEventListener('message', (event) => {
  // 防御式编程:始终校验数据
  try {
    const data = JSON.parse(event.data)
    console.log('收到消息:', data)
  } catch {
    console.warn('收到非 JSON 格式消息:', event.data)
  }
})

eventSource.addEventListener('notification', (event) => {
  // 自定义事件类型
  const data = JSON.parse(event.data)
  showNotification(data)
})

eventSource.addEventListener('error', (event) => {
  // EventSource 会自动重连,这里只做日志记录
  console.error('SSE 连接错误:', event)
})

SSE 协议格式

// 服务端返回的格式(Content-Type: text/event-stream)
data: {"type": "message", "content": "hello"}

data: {"type": "notification", "count": 5}

event: customEvent
data: {"key": "value"}

: 注释行,客户端忽略

每个 SSE 消息可以包含以下字段:

字段 必填 说明
data 消息内容,可以是多行
event 事件类型,对应 addEventListener 的事件名
id 消息 ID,用于断线重连时 Last-Event-ID
retry 重连间隔(ms)

3.2 自动重连机制

EventSource 内置了自动重连,这是它相比 WebSocket 最大的优势之一。

工作方式

  1. 连接断开时,浏览器自动尝试重连
  2. 默认重连间隔为 3 秒(可由服务端通过 retry 字段覆盖)
  3. 如果服务端发送了 id 字段,重连时会自动带上 Last-Event-ID 请求头,服务端可以从该 ID 之后继续推送
// 服务端通过 retry 字段控制重连间隔
// 服务端发送:retry: 5000
// 表示客户端 5 秒后重试

// 客户端断线后,EventSource 会自动重连
// 开发者不需要写任何重连代码

// 唯一需要处理的:服务端故障时,重连可能会无限循环
// 解决方案:封装超时断开机制
class SafeEventSource {
  private es: EventSource | null = null
  private maxRetries = 10
  private retryCount = 0
  private onMaxRetries: () => void

  constructor(url: string, onMaxRetries: () => void) {
    this.onMaxRetries = onMaxRetries
    this.connect(url)
  }

  private connect(url: string) {
    this.es = new EventSource(url)

    this.es.onopen = () => {
      this.retryCount = 0 // 连接成功,重置计数
    }

    this.es.onerror = () => {
      this.retryCount++
      if (this.retryCount > this.maxRetries) {
        console.error('[SSE] 重连次数过多,停止连接')
        this.es?.close()
        this.onMaxRetries()
      }
    }
  }

  close() {
    this.es?.close()
  }
}

3.3 与 ChatGPT 流式输出类似的场景

AI 流式输出是 SSE 的典型应用场景。用户发送提问后,服务端将 AI 生成的文本逐字/逐段推送给前端。

Client                          Server(AI 推理)
  |                               |
  |--- POST /api/chat             |
  |← event: chunk                 |
  |← data: {"text": "你好"}       |
  |← event: chunk                 |
  |← data: {"text": ","}         |
  |← event: chunk                 |
  |← data: {"text": "今天"}       |
  |← event: done                  |
  |← data: {"usage": {...}}       |
  |                               |
// src/composables/useSSEStream.ts — 流式输出封装

interface StreamOptions {
  url: string
  method?: 'GET' | 'POST'
  body?: Record<string, unknown>
  headers?: Record<string, string>
  onChunk: (text: string) => void
  onDone?: () => void
  onError?: (error: Error) => void
}

export function useSSEStream() {
  const abortController = ref<AbortController | null>(null)
  const isStreaming = ref(false)

  async function startStream(options: StreamOptions) {
    abortController.value = new AbortController()
    isStreaming.value = true

    try {
      const response = await fetch(options.url, {
        method: options.method ?? 'GET',
        headers: {
          'Content-Type': 'application/json',
          Accept: 'text/event-stream',
          ...options.headers,
        },
        body: options.body ? JSON.stringify(options.body) : undefined,
        signal: abortController.value.signal,
      })

      if (!response.ok) {
        throw new Error(`HTTP ${response.status}: ${response.statusText}`)
      }

      // 防御式编程:检查响应头确认是否为 SSE
      const contentType = response.headers.get('content-type') ?? ''
      if (!contentType.includes('text/event-stream')) {
        // 不是 SSE 格式时,按普通 JSON 读取
        const data = await response.json()
        options.onChunk(JSON.stringify(data))
        return
      }

      const reader = response.body?.getReader()
      if (!reader) {
        throw new Error('ReadableStream 不可用')
      }

      const decoder = new TextDecoder()
      let buffer = ''

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

        buffer += decoder.decode(value, { stream: true })

        // 解析 SSE 数据行
        const lines = buffer.split('\n')
        buffer = lines.pop() ?? '' // 最后一个可能是不完整行

        for (const line of lines) {
          if (line.startsWith('data: ')) {
            const dataStr = line.slice(6)

            // 防御式编程:忽略注释和空数据
            if (!dataStr || dataStr === '[DONE]') {
              if (dataStr === '[DONE]') {
                options.onDone?.()
              }
              continue
            }

            try {
              const data = JSON.parse(dataStr)
              if (data?.text !== undefined) {
                options.onChunk(data.text)
              } else if (data?.content !== undefined) {
                options.onChunk(data.content)
              }
            } catch {
              // 非 JSON 格式,直接输出原文
              options.onChunk(dataStr)
            }
          } else if (line.startsWith('event: ')) {
            // 处理自定义事件类型,可根据需要扩展
            const eventType = line.slice(7)
            if (eventType === 'done') {
              options.onDone?.()
            }
          }
        }
      }
    } catch (error: unknown) {
      if ((error as Error)?.name === 'AbortError') {
        // 主动取消,不触发 onError
        return
      }
      options.onError?.(error as Error)
    } finally {
      isStreaming.value = false
    }
  }

  function stopStream() {
    abortController.value?.abort()
    isStreaming.value = false
  }

  return { startStream, stopStream, isStreaming }
}
<!-- src/components/AIChat.vue — Vue 3 + AI 流式输出 -->
<script setup lang="ts">
import { ref } from 'vue'
import { ElInput, ElButton } from 'element-plus'
import { useSSEStream } from '@/composables/useSSEStream'

const { startStream, stopStream, isStreaming } = useSSEStream()
const prompt = ref('')
const output = ref('')

async function handleSubmit() {
  if (!prompt.value.trim()) return

  output.value = ''

  await startStream({
    url: '/api/chat/completions',
    method: 'POST',
    body: { prompt: prompt.value, stream: true },
    onChunk(text: string) {
      output.value += text
    },
    onDone() {
      console.log('流式输出完成')
    },
    onError(error) {
      ElMessage.error('请求失败: ' + error.message)
    },
  })
}

function handleStop() {
  stopStream()
}
</script>

<template>
  <div class="ai-chat">
    <ElInput
      v-model="prompt"
      type="textarea"
      :rows="3"
      placeholder="输入你的问题..."
      :disabled="isStreaming"
    />
    <div class="actions">
      <ElButton
        type="primary"
        @click="handleSubmit"
        :loading="isStreaming"
        :disabled="isStreaming"
      >
        发送
      </ElButton>
      <ElButton
        v-if="isStreaming"
        type="danger"
        @click="handleStop"
      >
        停止
      </ElButton>
    </div>
    <div class="output">
      <!-- 使用 v-html 时需要确保内容安全,防止 XSS -->
      <pre>{{ output }}</pre>
    </div>
  </div>
</template>
// src/components/AIChat.tsx — React + AI 流式输出

import { useState } from 'react'
import { Input, Button } from 'antd'
import { useSSEStream } from '@/hooks/useSSEStream'

const AIChat: React.FC = () => {
  const { startStream, stopStream, isStreaming } = useSSEStream()
  const [prompt, setPrompt] = useState('')
  const [output, setOutput] = useState('')

  const handleSubmit = async () => {
    if (!prompt.trim()) return
    setOutput('')

    await startStream({
      url: '/api/chat/completions',
      method: 'POST',
      body: { prompt, stream: true },
      onChunk(text: string) {
        setOutput((prev) => prev + text)
      },
      onError(error) {
        message.error('请求失败: ' + error.message)
      },
    })
  }

  return (
    <div className="ai-chat">
      <Input.TextArea
        value={prompt}
        onChange={(e) => setPrompt(e.target.value)}
        rows={3}
        placeholder="输入你的问题..."
        disabled={isStreaming}
      />
      <div className="actions">
        <Button
          type="primary"
          onClick={handleSubmit}
          loading={isStreaming}
          disabled={isStreaming}
        >
          发送
        </Button>
        {isStreaming && (
          <Button danger onClick={stopStream}>
            停止
          </Button>
        )}
      </div>
      <pre>{output}</pre>
    </div>
  )
}

export default AIChat

3.4 useSSE Composable / Hook 封装

// src/composables/useSSE.ts — Vue 3 Composable

import { ref, computed, onUnmounted } from 'vue'

interface SSEOptions {
  url: string
  withCredentials?: boolean
  maxRetries?: number
  onMessage?: (data: unknown) => void
  onError?: (event: Event) => void
  /** 事件名到处理函数的映射,如 { notification: (data) => {...} } */
  events?: Record<string, (data: unknown) => void>
}

export function useSSE(options: SSEOptions) {
  const isConnected = ref(false)
  const retryCount = ref(0)
  let eventSource: EventSource | null = null
  let reconnectTimer: ReturnType<typeof setTimeout> | null = null

  const maxRetries = options.maxRetries ?? Infinity

  function connect() {
    close()

    eventSource = new EventSource(options.url, {
      withCredentials: options.withCredentials ?? false,
    })

    eventSource.onopen = () => {
      isConnected.value = true
      retryCount.value = 0
    }

    eventSource.onmessage = (event) => {
      try {
        const data = JSON.parse(event.data)
        options.onMessage?.(data)
      } catch {
        options.onMessage?.(event.data)
      }
    }

    eventSource.onerror = (event) => {
      isConnected.value = false
      options.onError?.(event)

      // EventSource 内置自动重连
      // 但可能需要限制重连次数
      retryCount.value++

      if (retryCount.value > maxRetries) {
        console.error('[SSE] 超过最大重连次数,停止连接')
        eventSource?.close()
        return
      }
    }

    // 注册自定义事件
    if (options.events) {
      for (const [eventName, handler] of Object.entries(options.events)) {
        eventSource.addEventListener(eventName, (event: MessageEvent) => {
          try {
            handler(JSON.parse(event.data))
          } catch {
            handler(event.data)
          }
        })
      }
    }
  }

  function close() {
    eventSource?.close()
    eventSource = null
    isConnected.value = false
  }

  onUnmounted(() => {
    close()
  })

  return {
    isConnected,
    retryCount,
    connect,
    close,
  }
}

4. 消息处理

4.1 消息队列在前端的模拟

后端有消息队列(RabbitMQ / Kafka),前端同样需要"消息队列"来管理和分发消息,避免消息处理函数过于混乱。

// src/utils/message-queue.ts — 前端消息队列

type MessageHandler = (payload: Record<string, unknown>) => void

interface Message {
  type: string
  payload: Record<string, unknown>
  id?: string
  timestamp?: number
}

class MessageQueue {
  private handlers = new Map<string, Set<MessageHandler>>()
  private buffer: Message[] = []  // 连接未就绪时缓存消息
  private isFlushing = false
  private maxBufferSize = 100

  /** 注册消息处理器 */
  on(type: string, handler: MessageHandler) {
    if (!this.handlers.has(type)) {
      this.handlers.set(type, new Set())
    }
    this.handlers.get(type)!.add(handler)

    // 返回取消订阅函数
    return () => {
      this.handlers.get(type)?.delete(handler)
    }
  }

  /** 分发消息 */
  dispatch(message: Message) {
    // 防御式编程:校验消息格式
    if (!message?.type) {
      console.warn('[MessageQueue] 收到无效消息:', message)
      return
    }

    const handlers = this.handlers.get(message.type)
    if (!handlers || handlers.size === 0) {
      // 没有处理器时缓存消息,防止丢失
      this.bufferMessage(message)
      return
    }

    // 执行所有处理器
    handlers.forEach((handler) => {
      try {
        handler(message.payload ?? {})
      } catch (error) {
        console.error(`[MessageQueue] 消息处理器 ${message.type} 出错:`, error)
      }
    })
  }

  /** 缓存消息(连接断开时的消息丢失防护) */
  private bufferMessage(message: Message) {
    if (this.buffer.length >= this.maxBufferSize) {
      // 达到缓存上限时丢弃最早的消息
      this.buffer.shift()
    }
    this.buffer.push(message)
  }

  /** 注册处理器后,可以回放已缓存的消息 */
  flushBuffer() {
    if (this.isFlushing) return
    this.isFlushing = true

    while (this.buffer.length > 0) {
      const message = this.buffer.shift()
      if (message) {
        this.dispatch(message)
      }
    }

    this.isFlushing = false
  }

  /** 清理 */
  clear() {
    this.handlers.clear()
    this.buffer = []
  }
}

// 全局单例
export const messageQueue = new MessageQueue()
// 使用示例

// 1. 注册处理器(在组件初始化时)
messageQueue.on('chat:message', (payload) => {
  messages.value.push(payload as ChatMessage)
})

messageQueue.on('notification', (payload) => {
  showNotification(payload)
})

messageQueue.on('system:heartbeat', () => {
  // 更新心跳时间
})

// 2. WebSocket 接收消息后分发
socket.on('message', (raw) => {
  messageQueue.dispatch({
    type: raw.type,
    payload: raw.data,
    id: raw.id,
    timestamp: raw.timestamp,
  })
})

4.2 消息幂等处理(去重)

实时通信中,由于网络重试或服务端重推,客户端可能收到重复消息。幂等处理是必须的。

// src/utils/idempotent.ts — 消息去重

class IdempotentManager {
  private processedIds = new Set<string>()
  private maxCacheSize = 1000

  /** 检查消息是否已处理,未处理则标记 */
  isProcessed(messageId: string): boolean {
    if (this.processedIds.has(messageId)) {
      return true // 已处理,跳过
    }

    this.processedIds.add(messageId)
    this.evictOldest()

    return false
  }

  /** 清理过期的 ID */
  private evictOldest() {
    if (this.processedIds.size > this.maxCacheSize) {
      // 删除最早的 200 个 ID
      const iterator = this.processedIds.values()
      for (let i = 0; i < 200; i++) {
        const id = iterator.next()
        if (id.done) break
        this.processedIds.delete(id.value)
      }
    }
  }

  /** 清理所有记录 */
  clear() {
    this.processedIds.clear()
  }
}

export const idempotentManager = new IdempotentManager()
// 在消息处理器中使用

import { idempotentManager } from '@/utils/idempotent'

messageQueue.on('chat:message', (payload) => {
  const msg = payload as { id: string; content: string }

  // 防御式编程:如果服务端没有生成消息 ID,客户端自行生成
  const messageId = msg.id ?? `client_${Date.now()}_${Math.random()}`

  // 幂等检查
  if (idempotentManager.isProcessed(messageId)) {
    console.log('[消息] 重复消息,跳过:', messageId)
    return
  }

  // 处理消息...
  messages.value.push(msg)
})

4.3 消息顺序保证

WebSocket 虽然基于 TCP,保证传输层顺序,但以下场景可能导致消息顺序错乱:

  1. 消息重试:ACK 超时重发的消息可能晚于后面的消息到达
  2. 多设备同步:同一用户在不同设备上的操作,服务端的推送顺序可能与时间顺序不一致
  3. 服务端多实例:多个 WebSocket 服务实例可能导致消息到达顺序与发送顺序不一致
// src/utils/message-order.ts — 消息顺序保证

interface OrderedMessage {
  id: string
  sequenceId: number
  payload: Record<string, unknown>
}

class MessageOrderGuard {
  private expectedSeq = 0
  private buffer = new Map<number, OrderedMessage>()
  private maxBufferSize = 200

  /** 推送消息,返回是否按顺序处理 */
  push(message: OrderedMessage): boolean {
    // 防御式编程:序列号缺失时不做顺序保证
    if (message.sequenceId === undefined) {
      return true // 直接处理,不做顺序约束
    }

    if (message.sequenceId === this.expectedSeq) {
      // 恰好是期望的下一条
      this.expectedSeq++
      this.drainBuffer()
      return true
    }

    if (message.sequenceId > this.expectedSeq) {
      // 未来的消息,先缓存
      if (this.buffer.size < this.maxBufferSize) {
        this.buffer.set(message.sequenceId, message)
      }
      return false // 暂不处理
    }

    // 比 expectedSeq 小,可能是重复消息,跳过
    return false
  }

  /** 检查缓存中是否有可以按序处理的消息 */
  private drainBuffer() {
    while (this.buffer.has(this.expectedSeq)) {
      const msg = this.buffer.get(this.expectedSeq)!
      this.buffer.delete(this.expectedSeq)
      // 实际处理消息...
      this.expectedSeq++
    }
  }

  reset(seq: number) {
    this.expectedSeq = seq
    this.buffer.clear()
  }
}

5. 应用场景

5.1 即时通讯(IM)

能力 实现方案
消息收发 WebSocket(双向通信)
离线消息 HTTP 接口拉取 + WebSocket 实时推送
消息送达 ACK 机制(消息确认)
已读回执 专用消息类型 + 服务端广播
输入状态 节流后的 typing 事件推送
在线状态 连接状态管理 + 心跳检测
// IM 场景下推荐的消息格式

interface IMessage {
  id: string           // 消息 ID(幂等去重)
  type: 'text' | 'image' | 'file' | 'system'
  content: string
  senderId: number
  conversationId: string
  timestamp: number
  status: 'sending' | 'sent' | 'delivered' | 'read' | 'failed'
  /** 防御式编程:服务端可能扩展字段 */
  [key: string]: unknown
}

5.2 通知推送

方案 场景 实时性要求
SSE 页面内的通知提醒、消息中心未读数 秒级
WebSocket 需要交互的通知(审批、邀请) 秒级
Service Worker 浏览器关闭后的推送 依赖浏览器
推送平台 移动端推送(APNs / FCM) 分钟级

5.3 协作编辑(在线文档)

协作编辑的场景对实时通信有特殊要求:

挑战 解决方案
操作冲突 CRDT(Conflict-free Replicated Data Types)或 OT(Operational Transform)
光标同步 WebSocket 广播光标位置(需节流,每 50ms)
增量更新 WebSocket 传输操作指令而非完整文档
离线编辑 客户端缓存操作序列,重连后批量同步
// 协作编辑场景下的消息类型

// 操作指令(不传输完整文档)
interface Operation {
  type: 'insert' | 'delete' | 'replace'
  position: number
  length?: number
  text?: string
  userId: number
  version: number  // 文档版本号,用于冲突检测
}

// 光标同步(节流发送)
interface CursorSync {
  type: 'cursor'
  userId: number
  position: number
  selectionStart?: number
  selectionEnd?: number
  timestamp: number
}

5.4 数据大屏实时刷新

数据大屏场景下,SSE 是最推荐的方案,因为它是单向数据推送,实现简单且自带重连。

// src/composables/useRealtimeDashboard.ts

import { ref } from 'vue'
import { useSSE } from '@/composables/useSSE'

interface DashboardData {
  activeUsers: number
  todayOrders: number
  revenue: number
  cpuUsage: number
  memoryUsage: number
}

function useRealtimeDashboard(dashboardId: string) {
  const data = ref<DashboardData>({
    activeUsers: 0,
    todayOrders: 0,
    revenue: 0,
    cpuUsage: 0,
    memoryUsage: 0,
  })
  const alertMessage = ref('')

  const sse = useSSE({
    url: `/api/dashboard/${dashboardId}/stream`,
    onMessage(raw) {
      const payload = raw as Partial<DashboardData>

      // 防御式编程:只更新存在的字段,防止恶意数据注入
      if (typeof payload.activeUsers === 'number') {
        data.value.activeUsers = payload.activeUsers
      }
      if (typeof payload.todayOrders === 'number') {
        data.value.todayOrders = payload.todayOrders
      }
      // ...
    },
    events: {
      alert(data) {
        alertMessage.value = (data as { message: string }).message ?? ''
      },
    },
  })

  return { data, alertMessage, ...sse }
}

5.5 AI 流式输出

这是 SSE 在 AI 时代最典型的应用场景,在第 3.3 节已有详细示例。核心要点:

要点 说明
逐字推送 每推送一个 token,使用 data: { text: "x" }
停止生成 使用 AbortController 取消请求
错误处理 流中断时重试,但避免重复消费
连接复用 同一页面多个 AI 请求使用不同 URL 或请求 ID
UI 更新 不要使用 v-html 渲染未转义的 Markdown,使用专门的 Markdown 渲染组件
性能优化 大量文本输出时使用 requestAnimationFrame 控制 UI 更新频率

6. 注意事项与踩坑点

6.1 连接数限制

限制来源 限制数量 说明
浏览器 HTTP/1.1 同一域名 6~8 个 SSE 和 WebSocket 计数在此限制内(HTTP/1.1 时)
HTTP/2 同一域名 100 个 多路复用,推荐使用 H2
服务端连接数 取决于服务器配置 WebSocket 为长连接,需要关注内存和文件句柄
CDN / 反向代理 取决于供应商 某些 CDN 不支持 WebSocket

最佳实践

  • 使用 HTTP/2 减少连接数限制的影响
  • WebSocket 连接采用单例模式,全站共享一条连接
  • 页面关闭时主动释放连接(beforeunload 事件)
  • 如果在多 Tab 场景下有多个 WebSocket 连接,考虑使用 BroadcastChannel 在 Tab 间共享连接
// 多 Tab 共享 WebSocket 连接方案

// 主 Tab(实际建立 WebSocket 连接):
// 使用 Service Worker 或 SharedWorker 作为连接代理

// 替代方案:使用 BroadcastChannel 通信
const channel = new BroadcastChannel('websocket')

// Tab 1(持有 WebSocket 连接)
socket.onmessage = (event) => {
  channel.postMessage({ type: 'message', data: event.data })
}

// Tab 2(只监听 BroadcastChannel)
channel.onmessage = (event) => {
  if (event.data.type === 'message') {
    // 处理消息
  }
}

6.2 Nginx 配置

使用 Nginx 反向代理 WebSocket 和 SSE 时,需要额外配置:

# WebSocket 反向代理配置
upstream websocket {
    server 127.0.0.1:3000;
    # 保持长连接
    keepalive 64;
}

server {
    listen 80;
    server_name example.com;

    # WebSocket 路径
    location /ws/ {
        proxy_pass http://websocket;
        proxy_http_version 1.1;

        # WebSocket 必须的升级头
        proxy_set_header Upgrade $http_upgrade;
        proxy_set_header Connection "upgrade";

        # 超时设置(必须大于心跳间隔)
        proxy_read_timeout 3600s;    # 1 小时
        proxy_send_timeout 3600s;

        # 传递真实 IP
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;

        # 缓冲区设置(WebSocket 不需要缓冲)
        proxy_buffering off;
    }

    # SSE 路径
    location /api/stream/ {
        proxy_pass http://websocket;
        proxy_http_version 1.1;

        # SSE 配置
        proxy_set_header Connection "";
        proxy_buffering off;          # 关闭缓冲,确保实时推送
        proxy_cache off;              # 禁用缓存
        chunked_transfer_encoding on; # 启用分块传输

        # 长连接超时
        proxy_read_timeout 86400s;    # 24 小时
        proxy_send_timeout 86400s;
    }
}

Nginx 配置要点

配置项 WebSocket SSE 说明
proxy_http_version 1.1 1.1 都需要 HTTP/1.1
proxy_set_header Upgrade 需要 不需要 WebSocket 升级头部
proxy_set_header Connection "upgrade" "" SSE 不需要升级
proxy_buffering off off 必须关闭,否则无法实时
proxy_read_timeout 3600s+ 86400s+ 大于心跳间隔

6.3 移动端优化

问题 影响 解决方案
NAT 超时 运营商 NAT 映射有超时(通常 5 分钟) 缩短心跳间隔至 30s
网络切换 WiFi 与 4G/5G 切换导致连接断开 监听 online/offline 事件触发的重连
后台冻结 iOS Safari 或 Android 在后台冻结 JS 使用 Service Worker 或 Web Worker
省电模式 限制后台网络活动 降低心跳频率或使用被动通知
并发连接限制 移动浏览器连接数更少 全站共享一个 WebSocket 连接
// 移动端网络切换监听

function setupNetworkListener(onReconnect: () => void) {
  // 网络恢复时自动重连
  window.addEventListener('online', () => {
    console.log('[网络] 已恢复,开始重连')
    onReconnect()
  })

  window.addEventListener('office' /* 手动拼写正确为 offline */, () => {
    console.log('[网络] 已断开')
  })

  // iOS 页面可见性变化(从后台切回前台时检查连接)
  document.addEventListener('visibilitychange', () => {
    if (document.visibilityState === 'visible') {
      // 页面从后台切回前台,检查 WebSocket 连接状态
      if (socket?.readyState !== WebSocket.OPEN) {
        console.log('[WebSocket] 页面恢复,连接已断开,尝试重连')
        onReconnect()
      }
    }
  })
}

6.4 安全性

风险 说明 防范措施
未授权连接 任何人可以连接到 WebSocket 连接时携带 Token,服务端校验
消息注入 恶意消息包含 XSS 攻击代码 渲染前转义所有消息内容
消息伪造 客户端伪造其他用户的消息 服务端做发送者身份校验
重放攻击 拦截消息后重复发送 消息 ID 幂等 + 时间戳校验
连接耗尽 恶意客户端大量建立连接 服务端限制单用户连接数
// WebSocket 连接鉴权

// 方案一:URL Query 参数(不推荐,Token 可能被记录到日志)
const ws = new WebSocket(`wss://example.com/ws?token=${token}`)

// 方案二:协议头(推荐,但原生 WebSocket 不支持自定义 Header)
// 需要使用 Socket.IO 或 WebSocket 子协议

// 方案三:连接后立即鉴权(推荐)
// 1. 建立 WebSocket 连接
const ws = new WebSocket('wss://example.com/ws')

// 2. 连接建立后立即发送鉴权消息
ws.onopen = () => {
  ws.send(JSON.stringify({
    type: 'auth',
    token: 'your-jwt-token',
  }))
}

// 3. 服务端校验后决定是否保留连接
// 鉴权失败时服务端直接关闭连接
// 消息内容转义(防止 XSS)

function sanitizeMessage(text: string): string {
  const div = document.createElement('div')
  div.textContent = text // 自动转义 HTML
  return div.innerHTML
}

// Vue 中使用
// <div v-html="sanitizeMessage(msg.content)"></div>

// React 中使用
// <div dangerouslySetInnerHTML={{ __html: sanitizeMessage(msg.content) }} />

6.5 性能监控与调试

// WebSocket 性能监控

class WebSocketMonitor {
  private metrics = {
    messagesSent: 0,
    messagesReceived: 0,
    bytesSent: 0,
    bytesReceived: 0,
    reconnections: 0,
    lastPingLatency: 0,
    connectionDuration: 0,
  }

  private startTime = 0

  recordSend(data: string | ArrayBuffer) {
    this.metrics.messagesSent++
    this.metrics.bytesSent += typeof data === 'string'
      ? new Blob([data]).size
      : data.byteLength
  }

  recordReceive(data: string | ArrayBuffer) {
    this.metrics.messagesReceived++
    this.metrics.bytesReceived += typeof data === 'string'
      ? new Blob([data]).size
      : data.byteLength
  }

  recordReconnection() {
    this.metrics.reconnections++
  }

  getReport() {
    return { ...this.metrics, uptime: Date.now() - this.startTime }
  }
}

总结

场景 推荐方案 实现要点
即时通讯 WebSocket / Socket.IO 心跳 + 重连 + ACK + 幂等
通知推送 SSE 自动重连,最简实现
AI 流式输出 Fetch + ReadableStream 逐 token 推送,AbortController
数据大屏 SSE 单工推送,适合高频更新
协作编辑 WebSocket 操作指令同步 + CRDT/OT
移动端 WebSocket 心跳 30s,监听 online/visibility
降级方案 长轮询 兼容不支持 SSE/WebSocket 的旧浏览器

最终的核心原则:

  1. 选对协议:单向推送用 SSE,双向通信用 WebSocket,不要用轮询代替实时推送
  2. 心跳保活:WebSocket 必须有心跳机制(30s 间隔),防止 NAT 超时断开
  3. 指数退避:断线重连使用指数退避 + 随机抖动,防止重连风暴
  4. 幂等处理:任何消息都要做好去重,网络重试可能导致重复消息
  5. 防御式编程:永远不信任后端返回的消息格式,每个字段都要校验
  6. 主动释放:页面销毁时主动关闭连接,避免资源泄漏
  7. 连接共享:全站共享一条 WebSocket 连接,必要时使用 BroadcastChannel 跨 Tab 同步