在这里插入图片描述

每日一句正能量

“那些走过去的坎坷,终会成为生命里最坚硬的铠甲。”
回头看,那些曾让你哭的事,最后都变成了让你笑的资本。吃过的苦,都会化作护身的甲,护你今后岁岁平安。

摘要

系列导读:上一篇(第九十二篇)我们深入探讨了 网络状态监听 机制,包括 netAvailablenetLost 等事件的捕获与处理。本篇将在此基础上,构建一套完整的 WebSocket 实时通信方案,实现长连接的建立、心跳保活、智能重连,并与网络状态监听深度联动,打造生产级高可用的实时通信客户端。


一、前言:为什么需要 WebSocket?

在移动应用开发中,实时通信场景无处不在:即时消息推送、在线游戏同步、股票行情刷新、IoT 设备状态上报……传统的 HTTP 轮询方案存在明显的性能瓶颈——频繁的请求建立与断开带来巨大的开销,且消息延迟取决于轮询间隔。

WebSocket 协议通过一次 HTTP 握手升级,在客户端与服务器之间建立 全双工、持久化 的长连接,使得双方可以随时主动发送数据,极大降低了通信延迟和资源消耗。

在 HarmonyOS 中,@ohos.net.webSocket(API 9+)和 @kit.NetworkKit(API 12+)提供了原生的 WebSocket 支持。本文将围绕以下核心目标展开:

  1. 深度理解 WebSocket 握手流程与数据帧格式;
  2. 掌握 HarmonyOS WebSocket API 的核心用法;
  3. 封装 生产级客户端,支持心跳保活、指数退避重连;
  4. 联动 网络状态监听,实现智能连接管理;
  5. 保障 消息可靠投递,引入消息队列与持久化机制。

二、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 模块进行了显著增强:

  • 系统级心跳:通过 pingIntervalpongTimeout 配置,系统底层自动处理 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 事件会被触发。此时我们采取以下策略:

  1. 保持现有连接:如果当前连接仍然活跃,暂不主动断开;
  2. 调整心跳频率:根据新网络类型调整 heartbeatInterval
  3. 监控连接质量:若后续检测到心跳超时,再触发重连,新连接将自然路由到新网络。

这种"被动检测 + 主动适配"的策略避免了频繁断连带来的消息丢失。


七、消息队列与持久化:保障可靠投递

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,完整实现了一套 生产级高可用实时通信方案。核心要点回顾:

  1. 握手与状态机:理解 HTTP 升级握手和 CLOSED → OPEN 状态流转是编写健壮客户端的基础;
  2. 系统级心跳:善用 HarmonyOS 6+ 的 pingInterval / pongTimeout,减少应用层负担;
  3. 指数退避重连delay = min(base × 2ⁿ, maxDelay) × jitter,平衡恢复速度与服务器压力;
  4. 网络状态联动:结合 netAvailable / netLost 实现智能连接管理,网络恢复立即重连、丢失时暂停;
  5. 消息队列持久化:通过关系型数据库保障断网期间消息不丢失,恢复后自动补发;
  6. 资源管理:页面生命周期与连接生命周期绑定,避免内存泄漏和电量浪费。

转载自:https://blog.csdn.net/u014727709/article/details/163399193
欢迎 👍点赞✍评论⭐收藏,欢迎指正

Logo

作为“人工智能6S店”的官方数字引擎,为AI开发者与企业提供一个覆盖软硬件全栈、一站式门户。

更多推荐