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

Angular 16中Ngx-Mqtt与HiveMQ Broker连接偶发中断需刷新浏览器

Angular 16 + ngx-mqtt 连接中断后无法自动恢复,需硬刷新修复

我们的项目基于Angular 16,使用npm包 ngx-mqtt@^16.1.0 连接HiveMQ MQTT Broker。管理页面会在构造函数中自动触发连接:

constructor(
    private readonly mqttNotificationsService: MqttNotificationsService,
    private readonly config: AppConfigService,
) {
    if (this.config.getConfig('mqttMessagesEnabled')) {
        this.mqttNotificationsService.startConnection();
    }        
}

问题现象

连接会偶尔中断,此时通知图标显示0条消息,切换页面再返回也无法获取最新消息,唯一解决办法是硬刷新浏览器——刷新后客户端能重新连接Broker并获取最新消息。

相关代码实现

以下是MqttNotificationsService中核心方法的代码:

startConnection(): void {
    if (!this.config.getConfig('mqttMessagesEnabled')) {
        return;
    }

    this.topicsService
        .getTopics(-1, '')
        .pipe(
            tap((topics) => {
                this.topics = topics;
            }),
        )
        .subscribe({
            next: (_) => {
                this.connection = this.getBrokerConnection();
                try {
                    this.mqttService.connect(this.connection);
                    this.hubListener();
                } catch (error) {
                    console.log('mqtt.connect error', error);
                }
            },
        });
}

restartMqtt(): void {
    this.destroyConnection();
    timer(1000).pipe(
        tap((_) => {
            this.startConnection();
        }),
    ).subscribe();
}

destroyConnection() {
    this.mqttConnected = false;
    if (this.subscriptions) {
        this.subscriptions.forEach((sub) => sub.unsubscribe());
    }
    this.subscriptions = [];
    this.subscribeMultiple = [];
    try {
        this.mqttService?.disconnect(true);
    } catch (error) {
        console.log('Disconnect failed', error.toString());
    }
}

// ON CONNECT!
public hubListener = () => {
    this.mqttService.onConnect
        .pipe(
            tap((conn) => {
                console.log(`Mqtt OnConnect returned: ${conn.cmd}`);
                this.mqttConnected = true;
                this.subscribeToTopics();
            }),
        )
        .subscribe();
    this.mqttService.onClose.subscribe((_) => (this.mqttConnected = false));
    this.mqttService.onError
        .pipe(
            tap((result) => {
                console.log('Mqtt onError event fired: ', result);
            }),
        )
        .subscribe();
    this.mqttService.onOffline
        .pipe(
            tap((result) => {
                this.mqttConnected = false;
                console.log('Mqtt onOffline event fired: ', result);
            }),
        )
        .subscribe();
    this.mqttService.onMessage
        .subscribe();
    this.mqttService.onSuback
        .subscribe();
    this.mqttService.onPacketsend.subscribe();
};

问题分析与修复方案

核心问题

当前代码仅在连接断开时更新状态,未触发自动重连逻辑;且hubListener每次调用都会重复订阅事件,可能导致内存泄漏或逻辑冲突。

修复步骤

  1. 在断开事件中触发自动重连
    修改hubListener,在连接断开事件中调用重连逻辑,并新增事件订阅清理逻辑避免重复订阅:

    // 新增:定义事件订阅变量
    private connectSub?: Subscription;
    private closeSub?: Subscription;
    private errorSub?: Subscription;
    private offlineSub?: Subscription;
    private messageSub?: Subscription;
    private subackSub?: Subscription;
    private packetSendSub?: Subscription;
    
    public hubListener = () => {
        // 先清理之前的事件订阅
        this.clearEventSubscriptions();
    
        this.connectSub = this.mqttService.onConnect
            .pipe(
                tap((conn) => {
                    console.log(`Mqtt OnConnect returned: ${conn.cmd}`);
                    this.mqttConnected = true;
                    this.subscribeToTopics();
                }),
            )
            .subscribe();
    
        this.closeSub = this.mqttService.onClose.subscribe((_) => {
            this.mqttConnected = false;
            this.restartMqtt();
        });
    
        this.errorSub = this.mqttService.onError
            .pipe(
                tap((result) => {
                    console.log('Mqtt onError event fired: ', result);
                    this.mqttConnected = false;
                    this.restartMqtt();
                }),
            )
            .subscribe();
    
        this.offlineSub = this.mqttService.onOffline
            .pipe(
                tap((result) => {
                    this.mqttConnected = false;
                    console.log('Mqtt onOffline event fired: ', result);
                    this.restartMqtt();
                }),
            )
            .subscribe();
    
        this.messageSub = this.mqttService.onMessage.subscribe((message) => {
            // 此处补充通知计数更新逻辑
            // 示例:this.updateNotificationCount(message);
        });
    
        this.subackSub = this.mqttService.onSuback.subscribe();
        this.packetSendSub = this.mqttService.onPacketsend.subscribe();
    };
    
    // 新增:清理事件订阅
    private clearEventSubscriptions() {
        this.connectSub?.unsubscribe();
        this.closeSub?.unsubscribe();
        this.errorSub?.unsubscribe();
        this.offlineSub?.unsubscribe();
        this.messageSub?.unsubscribe();
        this.subackSub?.unsubscribe();
        this.packetSendSub?.unsubscribe();
    }
    

    同时在destroyConnection中调用clearEventSubscriptions:

    destroyConnection() {
        this.mqttConnected = false;
        this.clearEventSubscriptions(); // 新增
        if (this.subscriptions) {
            this.subscriptions.forEach((sub) => sub.unsubscribe());
        }
        this.subscriptions = [];
        this.subscribeMultiple = [];
        try {
            this.mqttService?.disconnect(true);
        } catch (error) {
            console.log('Disconnect failed', error.toString());
        }
    }
    
  2. 配置MQTT内置重连参数
    在getBrokerConnection中添加重连与会话持久化配置,增强连接稳定性:

    private getBrokerConnection(): IMqttServiceOptions {
        return {
            hostname: this.config.getConfig('mqttHost'),
            port: this.config.getConfig('mqttPort'),
            protocol: 'wss', // 根据实际环境调整
            reconnectPeriod: 3000, // 重连间隔3秒
            connectTimeout: 10000, // 连接超时10秒
            clean: false, // 保留会话,重连后恢复订阅
            clientId: `angular-client-${localStorage.getItem('userId')}`, // 固定客户端ID,确保会话一致性
        };
    }
    
  3. 页面切换时检查连接状态
    在组件的生命周期钩子中检查连接状态,未连接则触发重连:

    ngOnInit(): void {
        if (this.config.getConfig('mqttMessagesEnabled') && !this.mqttNotificationsService.mqttConnected) {
            this.mqttNotificationsService.restartMqtt();
        }
    }
    

内容的提问来源于stack exchange,提问作者bob.mazzo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 04:37:34