You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

MQTT.js主线程连接AWS IoT频繁断连重连,Worker线程稳定问题排查

主线程MQTT连接频繁断连排查分析

问题背景

我使用mqttjs通过WSS预签名URL连接AWS IoT MQTT服务,应用中存在两个MQTT连接:一个来自主线程,另一个来自Web Worker线程,二者由同一个MqttService类创建,逻辑完全一致。但主线程连接频繁断开重连,Worker线程连接却非常稳定,从未断开或需要重连。

已知排除项:

  • 重试5次失败后会终止并新建连接,请求量问题可忽略
  • 客户端ID随机生成,不会被其他客户端踢下线
  • 已排除连接超时和WiFi断开的情况

可能的原因分析

1. 主线程事件循环阻塞

主线程负责UI渲染、用户交互、JS执行等所有同步任务,如果存在耗时操作(比如大量DOM操作、复杂计算、同步网络请求),会导致MQTT客户端的心跳包(keepalive)无法及时发送,AWS IoT会判定连接超时并主动断开。而Web Worker运行在独立线程,不受主线程阻塞影响,心跳包能按时发送。

2. 浏览器后台标签页资源限制

多数浏览器会对后台标签页施加CPU/网络限制,比如降低定时器精度、限制WebSocket的发送频率。主线程所在的标签页如果处于后台,可能导致MQTT的keepalive心跳无法按时触发,引发连接断开;而Web Worker的线程不受后台标签页的这类限制,能维持稳定连接。

3. WebSocket连接的优先级差异

浏览器对主线程创建的WebSocket和Worker创建的WebSocket可能分配不同的网络优先级。主线程的WebSocket可能因页面其他网络请求(如API请求、静态资源加载)被抢占带宽,导致心跳包延迟或丢失;Worker的WebSocket则相对独立,受干扰更少。

4. isOnline判断的潜在问题

代码中isOnline方法依赖window.$nuxt.isOnline,这个状态可能存在误判:

private isOnline() {
  return typeof window !== 'undefined' && window?.$nuxt?.isOnline;
}

如果Nuxt的在线状态检测出现偏差,可能导致主线程的MQTT客户端被错误地调用disconnect;而Worker线程中没有window对象,isOnline始终返回false,不会触发这个逻辑,因此连接不会被主动断开。

5. transformWsUrl的客户端ID重置逻辑问题

在transformWsUrl回调中,主线程和Worker都会重置客户端ID,但主线程的回调执行时机可能和Worker存在差异:

transformWsUrl: (_url, _options, client) => {
  client.options.clientId = this.generateClientId();
  return this.presignedUrl;
}

如果主线程中该回调触发时,原连接尚未完全关闭,可能导致客户端ID冲突或连接参数异常,引发AWS IoT拒绝连接;Worker线程的执行环境更单纯,不存在这类时序问题。

6. 主线程内存波动与GC影响

主线程中频繁的DOM操作、对象创建会触发垃圾回收(GC),GC过程会暂停JS执行,可能导致MQTT客户端的心跳发送延迟。而Worker线程的内存波动小,GC频率低,对连接影响可忽略。

7. AWS IoT的连接配额限制

虽然客户端ID随机,但AWS IoT可能对同一IP或同一身份的并发连接存在隐性配额。主线程的频繁重连可能触发AWS的流量限制或连接频率限制,导致连接被主动断开;Worker连接稳定,不会触发这类限制。

排查建议

  • 监控主线程的事件循环延迟:使用performance.now()或浏览器DevTools的Performance面板,检查是否存在长任务阻塞。
  • 禁用后台标签页限制测试:将页面保持在前台,观察主线程连接是否恢复稳定。
  • 注释isOnline相关逻辑:临时修改disconnect方法,跳过isOnline判断,看主线程连接是否不再异常断开。
  • 增加MQTT日志:监听error事件,获取AWS IoT返回的断开原因(代码中未处理error事件,这是关键缺失):
this.client!.on('error', (err) => {
  logWithTimestamp(`[MQTT] [${this.caller}] error:`, err);
});
  • 调整keepalive时间:尝试将keepalive从15秒调整为30秒,减少心跳频率,降低主线程阻塞的影响。

附:MqttService类代码

/* eslint-disable no-useless-constructor, no-empty-function */

import mqtt, { MqttClient } from 'mqtt';
import { NuxtAxiosInstance } from '@nuxtjs/axios';
import RestService from './RestService';
import { Mqtt } from '~/types/Mqtt';
import { MILISECS_PER_SEC } from '~/configs';
import { logWithTimestamp } from '~/utils';

export type MqttServiceEventHandlers = {
  close?: Array<() => void>;
  disconnectd?: Array<() => void>;
  connected?: Array<() => void>;
  reconnect?: Array<() => void>;
  reconnected?: Array<() => void>;
  beforeReconect?: Array<() => void>;
};
export type MqttServiceEvent = keyof MqttServiceEventHandlers;
export interface IMqttService {
  httpClient?: NuxtAxiosInstance;
  presignedUrl?: string;
}

export class MqttService {
  public client: MqttClient | null = null;
  public closing = false;
  public reconnecting = false;
  // Example: "wss://abcdef123-ats.iot.us-east-2.amazonaws.com/mqtt?X-Amz-Algorithm=AWS4-HMAC-SHA256&X-Amz-Credential=XXXXXX%2F20230206%2Fus-east-2%2Fiotdevicegateway%2Faws4_request&X-Amz-Date=20230206T104907Z&X-Amz-Expires=900&X-Amz-SignedHeaders=host&X-Amz-Signature=abcxyz123"
  public presignedUrl = '';
  public httpClient = null as null | NuxtAxiosInstance;
  public retryCount = 0;
  public retryLimits = 5;
  public handlers: MqttServiceEventHandlers = {};

  constructor(
    { httpClient, presignedUrl }: IMqttService,
    public caller = 'main'
  ) {
    if (httpClient) {
      this.httpClient = httpClient;
    } else if (presignedUrl) {
      this.presignedUrl = presignedUrl;
    } else {
      throw new Error(
        '[MqttService] a httpClient or presigned URL must be provided'
      );
    }
  }

  async connect() {
    await this.updatePresignedUrl();
    this.client = mqtt.connect(this.presignedUrl, {
      clientId: this.generateClientId(),
      reconnectPeriod: 5000,
      connectTimeout: 30000,
      resubscribe: true,
      keepalive: 15,
      transformWsUrl: (_url, _options, client) => {
        // eslint-disable-next-line no-param-reassign
        client.options.clientId = this.generateClientId();
        logWithTimestamp(
          `[MQTT] [${this.caller}] transformWsUrl()`,
          client.options.clientId,
          this.signature
        );
        return this.presignedUrl;
      },
    });

    return this.setupHandlers();
  }

  protected setupHandlers() {
    return new Promise<MqttClient>((resolve, reject) => {
      this.client!.on('close', async () => {
        if (this.closing) return;

        if (this.retryCount >= this.retryLimits) {
          (this.handlers.close || []).forEach((handler) => handler());
          await this.disconnect();
          logWithTimestamp(`[MQTT] [${this.caller}] connection has closed!`);
          reject(new Error(`All retry attempts were failed (${this.caller})`));
          return;
        }

        if (this.retryCount === 0) {
          (this.handlers.beforeReconect || []).forEach((handler) => handler());
          logWithTimestamp(
            `[MQTT] [${this.caller}] connection lost`,
            this.presignedUrl
          );
        }

        // Re-generate new presigned URL at the 3rd attempt, or if the URL is expired
        if (this.retryCount === 2 || this.isExpired) {
          await this.updatePresignedUrl().catch(async () => {
            await this.disconnect();
            (this.handlers.close || []).forEach((handler) => handler());
            logWithTimestamp(
              `[MQTT] [${this.caller}] connection has closed! (Unable to get new presigned url)`
            );
          });
        }
      });

      this.client!.on('reconnect', () => {
        this.retryCount += 1;
        this.reconnecting = true;
        (this.handlers.reconnect || []).forEach((handler) => handler());
        logWithTimestamp(
          `[MQTT] [${this.caller}] reconnect (#${this.retryCount})`
        );
      });

      this.client!.on('connect', () => {
        if (this.reconnecting) {
          (this.handlers.reconnected || []).forEach((handler) => handler());
        }
        this.retryCount = 0;
        this.reconnecting = false;
        (this.handlers.connected || []).forEach((handler) => handler());
        logWithTimestamp(`[MQTT] [${this.caller}] connected`);
        resolve(this.client!);
      });
    });
  }

  disconnect({ force = true, silent = false, ...options } = {}) {
    this.closing = true;
    return new Promise<void>((resolve) => {
      if (this.client && this.isOnline()) {
        this.client.end(force, options, () => {
          this.reset(silent, '(fully)');
          resolve();
        });
      } else {
        this.client?.end(force);
        this.reset(silent, '(client-side only)');
        resolve();
      }
    });
  }

  reset(silent = false, debug?: any) {
    this.client = null;
    this.retryCount = 0;
    this.reconnecting = false;
    this.presignedUrl = '';
    this.closing = false;
    if (!silent) {
      (this.handlers.disconnectd || []).forEach((handler) => handler());
    }
    logWithTimestamp(`[MQTT] [${this.caller}] disconnected!`, {
      silent,
      debug,
    });
  }

  async destroy() {
    await this.disconnect({ silent: true });
    this.handlers = {};
  }

  // Get or assign a new wss url
  async updatePresignedUrl(url?: string) {
    if (this.httpClient) {
      const service = new RestService<Mqtt>(this.httpClient, '/mqtts');
      const { data } = await service.show('wss_link');
      this.presignedUrl = data!.wss_link;
    } else if (url) {
      this.presignedUrl = url;
    }
    return this.presignedUrl;
  }

  on(event: MqttServiceEvent, handler: () => void) {
    const { [event]: eventHanlders = [] } = this.handlers;
    eventHanlders.push(handler);
    this.handlers[event] = eventHanlders;
  }

  off(event: MqttServiceEvent, handler: () => void) {
    const { [event]: eventHanlders = [] } = this.handlers;
    const index = eventHanlders.findIndex((_handler) => _handler === handler);
    eventHanlders.splice(index, 1);
  }

  get date() {
    const matched = this.presignedUrl.match(/(X-Amz-Date=)(\w+)/);
    if (!matched) return null;

    return new Date(
      String(matched[2]).replace(
        /^(\d{4})(\d{2})(\d{2})(T\d{2})(\d{2})(\d{2}Z)$/,
        (__: string, YYYY: string, MM: string, DD: string, HH: string, mm: string, ss: string) => `${YYYY}-${MM}-${DD}${HH}:${mm}:${ss}`
      )
    );
  }

  get expires() {
    const matched = this.presignedUrl.match(/(X-Amz-Expires=)(\d+)/);
    return matched ? Number(matched[2]) : null;
  }

  get signature() {
    const matched = this.presignedUrl.match(/(X-Amz-Signature=)(\w+)/);
    return matched ? matched[2] : null;
  }

  get expiresDate() {
    if (!(this.date && this.expires)) return null;

    return new Date(this.date.getTime() + this.expires * MILISECS_PER_SEC);
  }

  get isExpired() {
    return !this.expiresDate || this.expiresDate <= new Date();
  }

  private generateClientId() {
    return `mqttjs_[${this.caller}]_${Math.random()
      .toString(16)
      .substring(2, 10)}`.toUpperCase();
  }

  private isOnline() {
    return typeof window !== 'undefined' && window?.$nuxt?.isOnline;
  }
}

内容的提问来源于stack exchange,提问作者Son Tr.

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 18:40:39