WebSocket 实时通信——从长连接到生产级高可用客户端
文章目录

每日一句正能量
“那些走过去的坎坷,终会成为生命里最坚硬的铠甲。”
回头看,那些曾让你哭的事,最后都变成了让你笑的资本。吃过的苦,都会化作护身的甲,护你今后岁岁平安。
摘要
系列导读:上一篇(第九十二篇)我们深入探讨了 网络状态监听 机制,包括
netAvailable、netLost等事件的捕获与处理。本篇将在此基础上,构建一套完整的 WebSocket 实时通信方案,实现长连接的建立、心跳保活、智能重连,并与网络状态监听深度联动,打造生产级高可用的实时通信客户端。
一、前言:为什么需要 WebSocket?
在移动应用开发中,实时通信场景无处不在:即时消息推送、在线游戏同步、股票行情刷新、IoT 设备状态上报……传统的 HTTP 轮询方案存在明显的性能瓶颈——频繁的请求建立与断开带来巨大的开销,且消息延迟取决于轮询间隔。
WebSocket 协议通过一次 HTTP 握手升级,在客户端与服务器之间建立 全双工、持久化 的长连接,使得双方可以随时主动发送数据,极大降低了通信延迟和资源消耗。
在 HarmonyOS 中,@ohos.net.webSocket(API 9+)和 @kit.NetworkKit(API 12+)提供了原生的 WebSocket 支持。本文将围绕以下核心目标展开:
- 深度理解 WebSocket 握手流程与数据帧格式;
- 掌握 HarmonyOS WebSocket API 的核心用法;
- 封装 生产级客户端,支持心跳保活、指数退避重连;
- 联动 网络状态监听,实现智能连接管理;
- 保障 消息可靠投递,引入消息队列与持久化机制。
二、WebSocket 核心原理
2.1 握手升级:从 HTTP 到 WebSocket
WebSocket 连接的建立始于一次特殊的 HTTP 请求。客户端发送带有 Upgrade: websocket 头部的 GET 请求,服务器返回 101 Switching Protocols 状态码,表示同意升级。此后,TCP 连接被"劫持"为 WebSocket 通道,双方开始基于 数据帧(Frame) 进行通信。

上图展示了完整的通信架构:HarmonyOS 客户端通过 @ohos.net.webSocket 模块发起连接,经 HTTP 握手升级后进入全双工数据传输阶段。客户端内部包含消息队列、心跳管理器、重连控制器和网络状态监听四大核心组件,共同保障连接的稳定性。
2.2 数据帧格式与状态机
WebSocket 消息被封装在帧中传输,关键字段包括:
| 字段 | 大小 | 说明 |
|---|---|---|
| FIN | 1 bit | 是否为最后一帧 |
| Opcode | 4 bits | 操作码(0x1 文本 / 0x2 二进制 / 0x8 关闭 / 0x9 Ping / 0xA Pong) |
| Mask | 1 bit | 客户端必须掩码 |
| Payload len | 7/7+16/7+64 bits | 负载长度 |
| Payload data | 变长 | 实际数据 |
从状态机视角看,WebSocket 连接经历 CLOSED → CONNECTING → OPEN → CLOSING → CLOSED 的完整生命周期:

- CLOSED:初始状态,连接未建立或已关闭;
- CONNECTING:调用
connect()后进入,等待握手完成; - OPEN:握手成功,可以收发消息;
- CLOSING:调用
close()后进入,等待关闭帧确认; - CLOSED:连接完全关闭,可触发重连逻辑。
理解这一状态机对编写健壮的客户端至关重要——任何异常(on("error"))或非正常关闭(on("close"),code ≠ 1000)都应触发重连机制。
三、HarmonyOS WebSocket API 详解
3.1 基础 API 使用
HarmonyOS 提供了简洁的 WebSocket API,核心方法如下:
import { webSocket } from '@kit.NetworkKit';
import { BusinessError } from '@kit.BasicServicesKit';
// 1. 创建 WebSocket 实例
let ws = webSocket.createWebSocket();
// 2. 连接服务器(支持自定义请求头和子协议)
await ws.connect('wss://api.example.com/ws', {
header: {
'Authorization': 'Bearer your_token',
'X-Client-Version': '1.0.0'
},
protocols: ['json'],
// HarmonyOS 6+ 支持系统级心跳配置
pingInterval: 30000, // 每 30 秒自动发送 Ping
pongTimeout: 5000 // 5 秒内未收到 Pong 则断开
});
// 3. 发送消息
await ws.send('Hello Server'); // 文本消息
await ws.send(new ArrayBuffer(1024)); // 二进制消息
// 4. 接收消息
ws.on('message', (err: BusinessError, value: string | ArrayBuffer) => {
if (err) {
console.error('接收错误:', err.message);
return;
}
if (typeof value === 'string') {
console.info('收到文本:', value);
} else {
console.info('收到二进制数据, 长度:', value.byteLength);
}
});
// 5. 监听连接关闭
ws.on('close', (err: BusinessError, value: webSocket.CloseResult) => {
console.info(`连接关闭: code=${value.code}, reason=${value.reason}`);
});
// 6. 监听错误
ws.on('error', (err: BusinessError) => {
console.error('WebSocket 错误:', err.message);
});
// 7. 主动关闭
await ws.close(1000, 'Client closing');
3.2 HarmonyOS 6 新增特性
从 API version 23 开始,WebSocket 模块进行了显著增强:
- 系统级心跳:通过
pingInterval和pongTimeout配置,系统底层自动处理 Ping/Pong 帧,无需应用层手动实现; - 子协议协商:
protocols数组允许客户端声明支持的子协议,服务器选择后通过getConnectionInfo()返回实际协商结果; - 服务端支持:
createWebSocketServer()允许在 HarmonyOS 设备上搭建 WebSocket 服务端,适用于局域网内设备直连场景。
四、实战:生产级 WebSocket 客户端封装
基础 API 虽然简洁,但直接用于生产环境存在明显不足:缺乏自动重连、心跳管理、消息队列、网络状态联动等能力。本节将封装一个 高可用的 WebSocket 客户端。
4.1 设计目标
| 能力 | 说明 |
|---|---|
| 自动重连 | 连接断开后按指数退避策略自动重连 |
| 心跳保活 | 定时发送 Ping 帧,检测连接活性 |
| 消息队列 | 断网期间消息入队,恢复后批量补发 |
| 网络联动 | 结合网络状态监听,智能控制重连时机 |
| 状态回调 | 对外暴露连接状态、重连进度等事件 |
4.2 核心实现代码
// websocket/WebSocketManager.ets
import { webSocket } from '@kit.NetworkKit';
import { connection } from '@kit.NetworkKit';
import { BusinessError } from '@kit.BasicServicesKit';
import { relationalStore } from '@kit.ArkData';
import { promptAction } from '@kit.ArkUI';
// ==================== 类型定义 ====================
interface WebSocketConfig {
url: string;
protocols?: string[];
headers?: Record<string, string>;
heartbeatInterval?: number; // 心跳间隔(ms),默认 30000
heartbeatTimeout?: number; // 心跳超时(ms),默认 5000
maxReconnectAttempts?: number; // 最大重连次数,默认 10
initialReconnectDelay?: number; // 初始重连延迟(ms),默认 1000
maxReconnectDelay?: number; // 最大重连延迟(ms),默认 30000
backoffFactor?: number; // 退避因子,默认 2
jitterFactor?: number; // 抖动因子(0-1),默认 0.3
enableQueue?: boolean; // 是否启用消息队列,默认 true
}
interface QueuedMessage {
id: string;
content: string | ArrayBuffer;
timestamp: number;
retryCount: number;
}
interface WebSocketState {
isConnected: boolean;
isReconnecting: boolean;
reconnectAttempts: number;
lastHeartbeatTime: number;
}
// ==================== 主类 ====================
export class WebSocketManager {
private ws: webSocket.WebSocket | null = null;
private config: WebSocketConfig;
private state: WebSocketState;
// 定时器
private heartbeatTimer: number = -1;
private heartbeatTimeoutTimer: number = -1;
private reconnectTimer: number = -1;
private stableTimer: number = -1;
// 消息队列
private messageQueue: QueuedMessage[] = [];
private db: relationalStore.RdbStore | null = null;
// 网络监听
private netConnection: connection.NetConnection | null = null;
private isNetworkAvailable: boolean = true;
// 回调
private onMessageCallback: ((data: string | ArrayBuffer) => void) | null = null;
private onStateChangeCallback: ((state: WebSocketState) => void) | null = null;
private onReconnectAttemptCallback: ((attempt: number, delay: number) => void) | null = null;
constructor(config: WebSocketConfig) {
this.config = {
heartbeatInterval: 30000,
heartbeatTimeout: 5000,
maxReconnectAttempts: 10,
initialReconnectDelay: 1000,
maxReconnectDelay: 30000,
backoffFactor: 2,
jitterFactor: 0.3,
enableQueue: true,
...config
};
this.state = {
isConnected: false,
isReconnecting: false,
reconnectAttempts: 0,
lastHeartbeatTime: 0
};
}
// ==================== 初始化 ====================
async initialize(): Promise<void> {
// 初始化消息队列数据库
if (this.config.enableQueue) {
await this.initMessageQueueDB();
}
// 注册网络状态监听
this.registerNetworkListener();
}
// ==================== 数据库初始化 ====================
private async initMessageQueueDB(): Promise<void> {
const STORE_CONFIG: relationalStore.StoreConfig = {
name: 'websocket_queue.db',
securityLevel: relationalStore.SecurityLevel.S1
};
this.db = await relationalStore.getRdbStore(getContext(), STORE_CONFIG);
// 创建消息队列表
const CREATE_TABLE_SQL = `
CREATE TABLE IF NOT EXISTS message_queue (
id TEXT PRIMARY KEY,
content TEXT NOT NULL,
timestamp INTEGER NOT NULL,
retry_count INTEGER DEFAULT 0
)
`;
await this.db.executeSql(CREATE_TABLE_SQL);
// 加载未发送消息
const resultSet = await this.db.querySql(
'SELECT * FROM message_queue ORDER BY timestamp ASC'
);
while (resultSet.goToNextRow()) {
this.messageQueue.push({
id: resultSet.getString(resultSet.getColumnIndex('id')),
content: resultSet.getString(resultSet.getColumnIndex('content')),
timestamp: resultSet.getLong(resultSet.getColumnIndex('timestamp')),
retryCount: resultSet.getLong(resultSet.getColumnIndex('retry_count'))
});
}
resultSet.close();
console.info(`[WebSocket] 加载 ${this.messageQueue.length} 条待发送消息`);
}
// ==================== 网络状态监听 ====================
private registerNetworkListener(): void {
this.netConnection = connection.createNetConnection();
// 网络可用时立即重连
this.netConnection.on('netAvailable', () => {
console.info('[WebSocket] 网络恢复可用');
this.isNetworkAvailable = true;
if (!this.state.isConnected && !this.state.isReconnecting) {
this.state.reconnectAttempts = 0; // 重置重试计数
this.scheduleReconnect(0); // 立即重连
}
});
// 网络丢失时暂停重连
this.netConnection.on('netLost', () => {
console.info('[WebSocket] 网络丢失');
this.isNetworkAvailable = false;
this.clearReconnectTimer();
});
// 网络能力变化时调整策略
this.netConnection.on('netCapabilitiesChange', (data) => {
const caps = data.netCap;
const hasWifi = caps.bearerTypes?.includes(connection.NetBearType.BEARER_WIFI);
const hasCellular = caps.bearerTypes?.includes(connection.NetBearType.BEARER_CELLULAR);
console.info(`[WebSocket] 网络类型变化: WiFi=${hasWifi}, Cellular=${hasCellular}`);
// 切换到更稳定的网络时,重置心跳间隔
if (hasWifi) {
this.config.heartbeatInterval = 30000; // WiFi 下标准心跳
} else if (hasCellular) {
this.config.heartbeatInterval = 20000; // 移动网络下更频繁心跳
}
});
this.netConnection.register();
}
// ==================== 连接管理 ====================
async connect(): Promise<boolean> {
if (this.state.isConnected) {
console.warn('[WebSocket] 已处于连接状态');
return true;
}
// 清理旧连接
await this.cleanup();
this.ws = webSocket.createWebSocket();
try {
const options: webSocket.WebSocketRequestOptions = {
header: this.config.headers || {},
protocols: this.config.protocols || []
};
// HarmonyOS 6+ 使用系统级心跳
if (this.config.heartbeatInterval) {
options.pingInterval = this.config.heartbeatInterval;
options.pongTimeout = this.config.heartbeatTimeout;
}
await this.ws.connect(this.config.url, options);
this.state.isConnected = true;
this.state.isReconnecting = false;
this.state.reconnectAttempts = 0;
this.state.lastHeartbeatTime = Date.now();
this.setupEventListeners();
this.startStableTimer();
this.notifyStateChange();
// 连接成功后补发队列消息
await this.flushMessageQueue();
console.info('[WebSocket] 连接成功');
return true;
} catch (error) {
console.error('[WebSocket] 连接失败:', (error as BusinessError).message);
this.handleDisconnect();
return false;
}
}
private setupEventListeners(): void {
if (!this.ws) return;
// 接收消息
this.ws.on('message', (err: BusinessError, value: string | ArrayBuffer) => {
if (err) {
console.error('[WebSocket] 消息接收错误:', err.message);
return;
}
// 更新心跳时间
this.state.lastHeartbeatTime = Date.now();
// 透传给业务层
this.onMessageCallback?.(value);
});
// 连接关闭
this.ws.on('close', (err: BusinessError, value: webSocket.CloseResult) => {
console.info(`[WebSocket] 连接关闭: code=${value.code}, reason=${value.reason}`);
this.handleDisconnect();
});
// 连接错误
this.ws.on('error', (err: BusinessError) => {
console.error('[WebSocket] 连接错误:', err.message);
this.handleDisconnect();
});
}
// ==================== 断开处理 ====================
private handleDisconnect(): void {
if (!this.state.isConnected && !this.state.isReconnecting) {
return; // 已经处理过了
}
this.state.isConnected = false;
this.stopStableTimer();
this.stopHeartbeat();
this.notifyStateChange();
// 触发重连(如果网络可用)
if (this.isNetworkAvailable) {
this.scheduleReconnect();
}
}
// ==================== 智能重连 ====================
private scheduleReconnect(immediateDelay?: number): void {
if (this.state.isReconnecting) return;
if (this.state.reconnectAttempts >= (this.config.maxReconnectAttempts || 10)) {
console.error('[WebSocket] 达到最大重连次数,停止重连');
promptAction.showToast({ message: '连接失败,请检查网络后重试' });
return;
}
this.state.isReconnecting = true;
this.state.reconnectAttempts++;
const delay = immediateDelay !== undefined
? immediateDelay
: this.calculateReconnectDelay();
console.info(`[WebSocket] 将在 ${delay}ms 后进行第 ${this.state.reconnectAttempts} 次重连`);
this.onReconnectAttemptCallback?.(this.state.reconnectAttempts, delay);
this.notifyStateChange();
this.reconnectTimer = setTimeout(() => {
this.connect();
}, delay);
}
/**
* 指数退避 + 随机抖动算法
* delay = min(base * 2^n, maxDelay) * (1 - jitter/2 + random * jitter)
*/
private calculateReconnectDelay(): number {
const base = this.config.initialReconnectDelay || 1000;
const maxDelay = this.config.maxReconnectDelay || 30000;
const factor = this.config.backoffFactor || 2;
const jitter = this.config.jitterFactor || 0.3;
const n = this.state.reconnectAttempts - 1;
// 指数退避
const baseDelay = Math.min(base * Math.pow(factor, n), maxDelay);
// 随机抖动 (0.85 ~ 1.15)
const jitterMultiplier = 1 - jitter / 2 + Math.random() * jitter;
return Math.floor(baseDelay * jitterMultiplier);
}
private clearReconnectTimer(): void {
if (this.reconnectTimer !== -1) {
clearTimeout(this.reconnectTimer);
this.reconnectTimer = -1;
}
}
// ==================== 心跳管理(应用层兜底) ====================
private startHeartbeat(): void {
// HarmonyOS 6+ 系统级心跳已启用,应用层作为兜底
if (this.heartbeatTimer !== -1) return;
this.heartbeatTimer = setInterval(() => {
if (!this.state.isConnected || !this.ws) return;
// 检查心跳超时
const elapsed = Date.now() - this.state.lastHeartbeatTime;
if (elapsed > (this.config.heartbeatInterval || 30000) * 2) {
console.warn('[WebSocket] 心跳超时,强制断开重连');
this.ws.close(1001, 'Heartbeat timeout');
return;
}
// 发送应用层心跳(业务级 Ping)
this.ws.send(JSON.stringify({ type: 'ping', timestamp: Date.now() }))
.catch(err => console.error('[WebSocket] 心跳发送失败:', err));
}, this.config.heartbeatInterval || 30000);
}
private stopHeartbeat(): void {
if (this.heartbeatTimer !== -1) {
clearInterval(this.heartbeatTimer);
this.heartbeatTimer = -1;
}
if (this.heartbeatTimeoutTimer !== -1) {
clearTimeout(this.heartbeatTimeoutTimer);
this.heartbeatTimeoutTimer = -1;
}
}
// ==================== 稳定计时器 ====================
private startStableTimer(): void {
this.stopStableTimer();
// 连接稳定 60 秒后重置重试计数
this.stableTimer = setTimeout(() => {
console.info('[WebSocket] 连接稳定,重置重试计数');
this.state.reconnectAttempts = 0;
}, 60000);
}
private stopStableTimer(): void {
if (this.stableTimer !== -1) {
clearTimeout(this.stableTimer);
this.stableTimer = -1;
}
}
// ==================== 消息发送与队列 ====================
async send(data: string | ArrayBuffer): Promise<boolean> {
// 如果已连接,直接发送
if (this.state.isConnected && this.ws) {
try {
await this.ws.send(data);
return true;
} catch (error) {
console.error('[WebSocket] 发送失败,转入队列:', error);
}
}
// 未连接时入队
if (this.config.enableQueue) {
await this.enqueueMessage(data);
}
return false;
}
private async enqueueMessage(data: string | ArrayBuffer): Promise<void> {
const message: QueuedMessage = {
id: `msg_${Date.now()}_${Math.random().toString(36).substr(2, 9)}`,
content: typeof data === 'string' ? data : '[Binary]',
timestamp: Date.now(),
retryCount: 0
};
this.messageQueue.push(message);
// 持久化到数据库
if (this.db) {
await this.db.executeSql(
'INSERT INTO message_queue (id, content, timestamp, retry_count) VALUES (?, ?, ?, ?)',
[message.id, message.content, message.timestamp, message.retryCount]
);
}
console.info(`[WebSocket] 消息已入队: ${message.id}`);
}
private async flushMessageQueue(): Promise<void> {
if (this.messageQueue.length === 0 || !this.ws) return;
console.info(`[WebSocket] 开始补发 ${this.messageQueue.length} 条队列消息`);
const failedMessages: QueuedMessage[] = [];
for (const msg of this.messageQueue) {
try {
await this.ws.send(msg.content);
// 发送成功,从数据库删除
if (this.db) {
await this.db.executeSql(
'DELETE FROM message_queue WHERE id = ?',
[msg.id]
);
}
} catch (error) {
msg.retryCount++;
if (msg.retryCount < 3) {
failedMessages.push(msg);
}
}
}
this.messageQueue = failedMessages;
console.info(`[WebSocket] 队列补发完成,剩余 ${failedMessages.length} 条`);
}
// ==================== 资源清理 ====================
private async cleanup(): Promise<void> {
this.clearReconnectTimer();
this.stopHeartbeat();
this.stopStableTimer();
if (this.ws) {
try {
await this.ws.close();
} catch (error) {
// 忽略关闭错误
}
this.ws = null;
}
}
async destroy(): Promise<void> {
await this.cleanup();
if (this.netConnection) {
this.netConnection.unregister();
this.netConnection = null;
}
if (this.db) {
await this.db.close();
this.db = null;
}
this.messageQueue = [];
}
// ==================== 回调注册 ====================
onMessage(callback: (data: string | ArrayBuffer) => void): void {
this.onMessageCallback = callback;
}
onStateChange(callback: (state: WebSocketState) => void): void {
this.onStateChangeCallback = callback;
}
onReconnectAttempt(callback: (attempt: number, delay: number) => void): void {
this.onReconnectAttemptCallback = callback;
}
private notifyStateChange(): void {
this.onStateChangeCallback?.({ ...this.state });
}
// ==================== 状态查询 ====================
getState(): WebSocketState {
return { ...this.state };
}
getQueueLength(): number {
return this.messageQueue.length;
}
}
4.3 在 UI 中使用
// pages/ChatPage.ets
import { WebSocketManager } from '../websocket/WebSocketManager';
@Entry
@Component
struct ChatPage {
@State messages: string[] = [];
@State inputText: string = '';
@State connectionStatus: string = '未连接';
@State reconnectInfo: string = '';
@State queueCount: number = 0;
private wsManager: WebSocketManager | null = null;
private scroller: Scroller = new Scroller();
async aboutToAppear() {
this.wsManager = new WebSocketManager({
url: 'wss://your-server.com/ws',
headers: {
'Authorization': 'Bearer your_token'
},
heartbeatInterval: 25000,
maxReconnectAttempts: 10,
enableQueue: true
});
// 注册消息回调
this.wsManager.onMessage((data) => {
if (typeof data === 'string') {
this.messages.push(`[收] ${data}`);
this.scroller.scrollEdge(Edge.Bottom);
}
});
// 注册状态回调
this.wsManager.onStateChange((state) => {
if (state.isConnected) {
this.connectionStatus = '已连接';
} else if (state.isReconnecting) {
this.connectionStatus = `重连中(${state.reconnectAttempts})`;
} else {
this.connectionStatus = '未连接';
}
this.queueCount = this.wsManager?.getQueueLength() || 0;
});
// 注册重连进度回调
this.wsManager.onReconnectAttempt((attempt, delay) => {
this.reconnectInfo = `第${attempt}次重连,${delay}ms后尝试`;
});
// 初始化并连接
await this.wsManager.initialize();
await this.wsManager.connect();
}
aboutToDisappear() {
this.wsManager?.destroy();
}
async sendMessage() {
if (!this.inputText.trim()) return;
const success = await this.wsManager?.send(this.inputText);
if (success) {
this.messages.push(`[发] ${this.inputText}`);
} else {
this.messages.push(`[队列] ${this.inputText}`);
}
this.inputText = '';
this.scroller.scrollEdge(Edge.Bottom);
}
build() {
Column() {
// 状态栏
Row() {
Text('WebSocket 实时通信')
.fontSize(18)
.fontWeight(FontWeight.Bold)
.layoutWeight(1)
Column() {
Text(this.connectionStatus)
.fontSize(12)
.fontColor(this.connectionStatus.includes('已连接') ? '#7ED321' : '#F5A623')
Text(this.reconnectInfo)
.fontSize(10)
.fontColor('#999')
Text(`队列: ${this.queueCount}`)
.fontSize(10)
.fontColor('#666')
}
.alignItems(HorizontalAlign.End)
}
.width('100%')
.padding(15)
.backgroundColor('#F5F5F5')
// 消息列表
List({ scroller: this.scroller }) {
ForEach(this.messages, (msg: string, index: number) => {
ListItem() {
Text(msg)
.fontSize(14)
.padding(10)
.backgroundColor(msg.startsWith('[发]') ? '#E3F2FD' : '#F5F5F5')
.borderRadius(8)
.width('100%')
}
.padding({ top: 5, bottom: 5 })
}, (msg: string, index: number) => index.toString())
}
.width('100%')
.layoutWeight(1)
.padding(10)
// 输入区
Row() {
TextInput({ text: this.inputText, placeholder: '输入消息...' })
.layoutWeight(1)
.height(45)
.onChange((value) => this.inputText = value)
.onSubmit(() => this.sendMessage())
Button('发送')
.width(70)
.height(45)
.margin({ left: 10 })
.onClick(() => this.sendMessage())
}
.width('100%')
.padding(15)
.backgroundColor('#FFFFFF')
}
.width('100%')
.height('100%')
}
}
五、心跳保活与指数退避重连
5.1 心跳机制的双重保障
在生产环境中,WebSocket 连接可能因以下原因被静默断开:
- NAT 超时:运营商网络对空闲连接的回收时间通常为 5~15 分钟;
- 防火墙策略:企业防火墙可能主动断开长时间无数据的连接;
- 服务器重启:服务端升级或重启导致连接异常关闭。
HarmonyOS 6+ 提供了 系统级心跳(pingInterval / pongTimeout),由底层自动发送协议级 Ping/Pong 帧。同时,我们在应用层实现了业务级心跳作为兜底,形成双重保障。

上图展示了完整的时序:正常通信阶段客户端定期发送 Ping 帧并接收 Pong 响应;一旦心跳超时或收到 netLost 事件,连接被标记为断开,进入重连阶段。重连采用 指数退避算法,延迟时间依次为 1s、2s、4s、8s、16s……上限 30s,并引入 30% 的随机抖动避免"惊群效应"。当 netAvailable 事件触发时,立即跳过等待进行重连。
5.2 指数退避算法详解
private calculateReconnectDelay(): number {
const base = 1000; // 基础延迟 1s
const maxDelay = 30000; // 最大延迟 30s
const factor = 2; // 退避因子
const jitter = 0.3; // 抖动因子 30%
const n = this.reconnectAttempts - 1;
// 指数退避: 1s, 2s, 4s, 8s, 16s, 30s, 30s...
const baseDelay = Math.min(base * Math.pow(factor, n), maxDelay);
// 随机抖动: 实际延迟在 0.85~1.15 倍之间波动
const jitterMultiplier = 1 - jitter/2 + Math.random() * jitter;
return Math.floor(baseDelay * jitterMultiplier);
}
该算法的优势在于:
- 初期快速恢复:前几次重连间隔短,用户体验好;
- 避免服务器压力:后期间隔指数增长,防止故障期间大量客户端同时重连压垮服务器;
- 随机打散:抖动因子将重连请求均匀分布到时间轴上。
六、与网络状态监听的深度联动
上一篇文章我们实现了网络状态的实时监听。将其与 WebSocket 结合,可以实现更智能的连接管理:

6.1 联动策略矩阵
| 网络事件 | 触发动作 | 说明 |
|---|---|---|
netAvailable |
立即重连 | 网络恢复后无需等待退避计时器 |
netLost |
暂停重连 | 避免在无网络环境下无意义重试,节省电量 |
netCapabilitiesChange |
调整心跳间隔 | WiFi 下 30s,移动网络下 20s,弱网 15s |
netBlockStatusChange |
标记网络阻塞 | 阻塞期间消息强制入队 |
6.2 网络切换时的平滑迁移
当设备从 WiFi 切换到移动数据(或反之),netCapabilitiesChange 事件会被触发。此时我们采取以下策略:
- 保持现有连接:如果当前连接仍然活跃,暂不主动断开;
- 调整心跳频率:根据新网络类型调整
heartbeatInterval; - 监控连接质量:若后续检测到心跳超时,再触发重连,新连接将自然路由到新网络。
这种"被动检测 + 主动适配"的策略避免了频繁断连带来的消息丢失。
七、消息队列与持久化:保障可靠投递
7.1 为什么需要消息队列?
在弱网或断网场景下,用户可能持续发送消息。如果没有队列机制,这些消息将直接丢失。通过引入 内存队列 + 关系型数据库持久化 的双层保障,可以确保:
- 消息在发送前先入队并写入数据库;
- 连接恢复后按顺序批量补发;
- 应用重启后从数据库恢复未发送消息。
7.2 消息生命周期
业务调用 send()
↓
检查连接状态
├─ 已连接 → 直接发送 → 成功则结束
└─ 未连接 → 入队 + 写入数据库
↓
连接恢复后
↓
从队列取出 → 发送 → 成功则从数据库删除
↓
发送失败 → retryCount++ → 超过3次则丢弃并告警
7.3 数据库表结构
CREATE TABLE message_queue (
id TEXT PRIMARY KEY, -- 消息唯一标识
content TEXT NOT NULL, -- 消息内容(二进制消息标记为 [Binary])
timestamp INTEGER NOT NULL, -- 入队时间戳
retry_count INTEGER DEFAULT 0 -- 重试次数
);
八、性能优化与最佳实践
8.1 二进制数据传输优化
对于图片、语音、视频等大体积数据,建议使用 二进制帧 传输,并配合分片策略:
// 分片发送大文件
async function sendLargeFile(fileBuffer: ArrayBuffer, chunkSize: number = 64 * 1024) {
const totalChunks = Math.ceil(fileBuffer.byteLength / chunkSize);
for (let i = 0; i < totalChunks; i++) {
const start = i * chunkSize;
const end = Math.min(start + chunkSize, fileBuffer.byteLength);
const chunk = fileBuffer.slice(start, end);
const metadata = JSON.stringify({
type: 'file_chunk',
index: i,
total: totalChunks,
size: end - start
});
// 先发送元数据(文本帧)
await wsManager.send(metadata);
// 再发送数据(二进制帧)
await wsManager.send(chunk);
}
}
8.2 弱网环境下的自适应策略
- 自适应心跳:根据 RTT(往返时延)动态调整心跳间隔,RTT 高时缩短间隔;
- 消息压缩:对文本消息启用 Gzip 压缩后再发送;
- 优先级队列:将消息分为高/中/低优先级,弱网时优先发送高优先级消息。
8.3 内存与电量优化
- 及时清理:页面销毁时调用
destroy(),释放 WebSocket、定时器、数据库和网络监听资源; - 后台策略:应用进入后台时,可适当延长心跳间隔或暂停非关键消息发送;
- 连接复用:同一域名下的多个业务模块共享一个 WebSocket 连接,通过消息路由分发。
九、总结
本文从 WebSocket 协议原理出发,结合 HarmonyOS 原生 API,完整实现了一套 生产级高可用实时通信方案。核心要点回顾:
- 握手与状态机:理解 HTTP 升级握手和 CLOSED → OPEN 状态流转是编写健壮客户端的基础;
- 系统级心跳:善用 HarmonyOS 6+ 的
pingInterval/pongTimeout,减少应用层负担; - 指数退避重连:
delay = min(base × 2ⁿ, maxDelay) × jitter,平衡恢复速度与服务器压力; - 网络状态联动:结合
netAvailable/netLost实现智能连接管理,网络恢复立即重连、丢失时暂停; - 消息队列持久化:通过关系型数据库保障断网期间消息不丢失,恢复后自动补发;
- 资源管理:页面生命周期与连接生命周期绑定,避免内存泄漏和电量浪费。
转载自:https://blog.csdn.net/u014727709/article/details/163399193
欢迎 👍点赞✍评论⭐收藏,欢迎指正
更多推荐




所有评论(0)