实时通信实践(SSE + WebSocket)
在现代 Web 应用中,实时通信已从"加分项"变成了"基础能力"。无论是即时消息、协作编辑、数据大屏,还是 AI 流式输出,都离不开实时通信技术。
本文将系统性地讨论四种主流的实时通信方案(短轮询、长轮询、SSE、WebSocket),并重点展开 WebSocket 和 SSE 的最佳实践。
1. 实时通信方案对比
1.1 短轮询(Short Polling)
原理:前端通过 setInterval 或 setTimeout 定时向服务端发起 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 最大的优势之一。
工作方式:
- 连接断开时,浏览器自动尝试重连
- 默认重连间隔为 3 秒(可由服务端通过
retry字段覆盖) - 如果服务端发送了
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,保证传输层顺序,但以下场景可能导致消息顺序错乱:
- 消息重试:ACK 超时重发的消息可能晚于后面的消息到达
- 多设备同步:同一用户在不同设备上的操作,服务端的推送顺序可能与时间顺序不一致
- 服务端多实例:多个 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 的旧浏览器 |
最终的核心原则:
- 选对协议:单向推送用 SSE,双向通信用 WebSocket,不要用轮询代替实时推送
- 心跳保活:WebSocket 必须有心跳机制(30s 间隔),防止 NAT 超时断开
- 指数退避:断线重连使用指数退避 + 随机抖动,防止重连风暴
- 幂等处理:任何消息都要做好去重,网络重试可能导致重复消息
- 防御式编程:永远不信任后端返回的消息格式,每个字段都要校验
- 主动释放:页面销毁时主动关闭连接,避免资源泄漏
- 连接共享:全站共享一条 WebSocket 连接,必要时使用 BroadcastChannel 跨 Tab 同步