BehaviorSubject接收消息重复且数量指数增长的原因及解决方案咨询
我一眼就看到问题所在了——你这是典型的重复订阅BehaviorSubject + 重复注册SignalR事件监听导致的指数级消息爆炸,咱们一步步拆解:
问题根源分析
1. 重复注册SignalR事件监听
你的组件代码里,每次hubMessage$收到消息时,都会调用this.signalRService.newLocationRecieved(),而这个方法里的this.hubConnection.on('NewLoc', ...)是添加监听(不是替换)。每调用一次,就会多注册一个NewLoc事件的回调。
举个例子:
- 第一次收到消息,调用
newLocationRecieved(),注册1个监听; - 第二次服务端发消息,这个监听触发,又调用
newLocationRecieved(),注册第2个监听; - 第三次服务端发消息,2个监听都会触发,每个都注册新的监听,变成4个;
- 以此类推,消息数量自然呈指数翻倍。
2. 重复订阅BehaviorSubject
你的newCoordinate()方法每次被createMarkers()调用(而createMarkers()又被ngOnChanges()触发),都会对hubMessage$新增一个订阅。这些订阅不会自动销毁,导致每一条消息过来,所有订阅都会执行,进一步叠加消息数量。
解决方案
第一步:停止重复注册SignalR监听
把newLocationRecieved()的调用从组件的订阅回调中彻底移除,改成在SignalR连接建立时只注册一次监听。
修改你的SignalR服务代码:
hubMessage$ = new BehaviorSubject({}); private isConnected = false; // 新增状态变量,防止重复连接 public startConnection = (id: number) => { if (this.isConnected) { console.log('已建立连接,无需重复启动'); return; } this.hubConnection = new signalR.HubConnectionBuilder() .withUrl('https://api/hub') .build(); // 提前注册监听,只注册一次 this.hubConnection.on('NewLoc', (data) => { console.log('new location recieved', data); this.hubMessage$.next(data); }); this.hubConnection .start() .then(() => { console.log('connection established'); this.isConnected = true; this.sendDriverId(id); // 连接成功后再发送订阅请求 }) .catch(err => { console.log('Error while starting connection: ' + err); this.isConnected = false; this.retryConnection(); }); } // 如果需要更新监听,可以保留这个方法,但要先移除旧监听再添加 public newLocationRecieved() { this.hubConnection.off('NewLoc'); // 先移除旧监听 this.hubConnection.on('NewLoc', (data) => { console.log('new location recieved', data); this.hubMessage$.next(data); }); } public sendDriverId(id: number = 1) { this.hubConnection.send('SubOnDriver', { driverId: id }); }
第二步:管理组件订阅,避免重复订阅
在Angular中,订阅Observable后必须及时销毁,防止内存泄漏和重复执行。推荐两种方式:
方式一:手动管理Subscription
import { Subscription } from 'rxjs'; // 组件类中声明订阅变量 private hubSub?: Subscription; ngOnChanges() { this.dispatchDetails; this.createMarkers(); } createMarkers() { console.log('Connection start right now ', this.dispatchDetails); // 先取消之前的订阅,再重新订阅 this.hubSub?.unsubscribe(); this.setupCoordinateListener(); } private setupCoordinateListener() { this.hubSub = this.signalRService.hubMessage$.subscribe( (data: any) => { console.log('recieved new coordinate ?', data); // 这里删掉this.signalRService.newLocationRecieved()!! this.locationCoords = data; if (this.locationCoords.location) { this.latitude = this.locationCoords?.location?.latitude; this.longitude = this.locationCoords?.location?.longitude; } } ); } // 组件销毁时务必取消订阅 ngOnDestroy() { this.hubSub?.unsubscribe(); }
方式二:用async管道自动管理(更推荐)
这种方式不需要手动写subscribe,Angular会自动帮你管理订阅的创建和销毁:
组件类代码:
// 把BehaviorSubject暴露给模板 locationUpdates$ = this.signalRService.hubMessage$.pipe( tap((data: any) => { console.log('recieved new coordinate ?', data); this.locationCoords = data; if (this.locationCoords.location) { this.latitude = this.locationCoords?.location?.latitude; this.longitude = this.locationCoords?.location?.longitude; } }) ); ngOnChanges() { this.dispatchDetails; this.createMarkers(); } createMarkers() { console.log('Connection start right now ', this.dispatchDetails); // 这里不需要再手动调用订阅方法了 }
模板中(添加async管道):
<!-- 用ngIf结合async管道,确保只有当有数据时才处理 --> <div *ngIf="locationUpdates$ | async"></div>
第三步:确保SignalR连接只建立一次
通过服务中的isConnected状态变量,避免重复调用startConnection创建多个连接,这也能减少不必要的消息重复。
按照上面的修改,你的消息重复问题应该就能彻底解决了,不会再出现指数级增长的情况。
内容的提问来源于stack exchange,提问作者Segma

