Angular中WebSocketSubject无法正常断开连接问题求助
问题
我有一个Angular项目,需要通过WebSocket接收服务器数据。参考angular-websocket-starter实现了WebSocket服务,但该实现没有断开Socket的方法。我需要在调用Socket的组件销毁时关闭Socket,于是自行添加了disconnect方法,但组件销毁后Socket会自动重连并继续工作。以下是相关代码,请问需要修改服务中的哪些内容才能实现组件销毁时正常断开Socket?
websocket.service.ts
import { Injectable, OnDestroy, Inject } from '@angular/core'; import { Observable, SubscriptionLike, Subject, Observer, interval } from 'rxjs'; import { filter, map } from 'rxjs/operators'; import { WebSocketSubject, WebSocketSubjectConfig } from 'rxjs/webSocket'; import { share, distinctUntilChanged, takeWhile } from 'rxjs/operators'; import { IWebsocketService, IWsMessage, WebSocketConfig } from './websocket.interfaces'; import { config } from './websocket.config'; export const wsUrl = 'wss://example.com:1234'; @Injectable({ providedIn: 'root' }) export class WebsocketService implements IWebsocketService, OnDestroy { private config: WebSocketSubjectConfig<IWsMessage<any>>; private websocketSub: SubscriptionLike; private statusSub: SubscriptionLike; private reconnection$: Observable<number>; private websocket$: WebSocketSubject<IWsMessage<any>>; private connection$: Observer<boolean>; private wsMessages$: Subject<IWsMessage<any>>; private reconnectInterval: number; private reconnectAttempts: number; private isConnected: boolean; public status: Observable<boolean>; constructor(@Inject(config) private wsConfig: WebSocketConfig) { this.wsMessages$ = new Subject<IWsMessage<any>>(); this.reconnectInterval = wsConfig.reconnectInterval || 5000; // pause between connections this.reconnectAttempts = wsConfig.reconnectAttempts || 10; // number of connection attempts this.config = { url: wsConfig.url, closeObserver: { next: (event: CloseEvent) => { console.log('WebSocket disconnected!'); this.websocket$ = null; this.connection$.next(false); } }, openObserver: { next: (event: Event) => { console.log('socket event:',event); console.log('WebSocket connected!'); this.connection$.next(true); } } }; // connection status this.status = new Observable<boolean>((observer) => { this.connection$ = observer; }).pipe(share(), distinctUntilChanged()); // run reconnect if not connection this.statusSub = this.status .subscribe((isConnected) => { this.isConnected = isConnected; if (!this.reconnection$ && typeof(isConnected) === 'boolean' && !isConnected) { this.reconnect(); } }); this.websocketSub = this.wsMessages$.subscribe( null, (error: ErrorEvent) => console.error('WebSocket error!', error) ); this.connect(); } disconnect(): void { // <----我添加的方法 this.websocket$ = null; this.wsMessages$.complete(); this.connection$.next(false); this.connection$.complete(); } ngOnDestroy(): void { this.websocketSub.unsubscribe(); this.statusSub.unsubscribe(); } /* * 连接WebSocket * */ private connect(): void { this.websocket$ = new WebSocketSubject(this.config); this.websocket$.subscribe( (message) => this.wsMessages$.next(message), (error: Event) => { if (!this.websocket$) { // 出错时重连 this.reconnect(); } }); } /* * 未连接或出错时重连 * */ private reconnect(): void { this.reconnection$ = interval(this.reconnectInterval) .pipe(takeWhile((v, index) => index < this.reconnectAttempts && !this.websocket$)); this.reconnection$.subscribe( () => this.connect(), null, () => { // 重连尝试结束后完成Subject this.reconnection$ = null; if (!this.websocket$) { this.wsMessages$.complete(); this.connection$.complete(); } }); } /* * 监听消息事件 * */ public on<T>(event: string): Observable<T> { if (event) { return this.wsMessages$.pipe( filter((message: IWsMessage<T>) => message.event === event), map((message: IWsMessage<T>) => message.data) ); } } }
调用Socket的组件代码
ngOnInit(): void { this.websocketService.on<any[]>('messages').subscribe(res => { console.log('服务器响应:',res); }); }); ngOnDestroy(): void { this.websocketService.disconnect(); <---这个方法不起作用 }
解决方案
问题核心是你添加的disconnect方法既没有实际关闭WebSocket连接,也没有阻止自动重连机制。以下是具体修改步骤:
1. 添加状态标记与订阅变量
在服务类中新增两个私有变量,分别标记是否为主动断开,以及保存重连的订阅对象:
private isManualDisconnect = false; private reconnectionSubscription: SubscriptionLike; // 用于取消重连订阅
2. 完善disconnect方法
修改disconnect方法,确保关闭连接、终止数据流并阻止后续重连:
disconnect(): void { // 标记为主动断开,禁止后续自动重连 this.isManualDisconnect = true; // 关闭WebSocket连接 if (this.websocket$) { this.websocket$.complete(); // 调用complete会触发closeObserver,正确关闭连接 this.websocket$ = null; } // 终止消息流和状态流 this.wsMessages$.complete(); if (this.connection$) { this.connection$.next(false); this.connection$.complete(); } // 取消当前的重连订阅(如果存在) if (this.reconnectionSubscription) { this.reconnectionSubscription.unsubscribe(); this.reconnectionSubscription = null; this.reconnection$ = null; } }
3. 修改重连触发逻辑
状态订阅中的重连判断
在statusSub的订阅逻辑中,只有非主动断开时才触发重连:
this.statusSub = this.status .subscribe((isConnected) => { this.isConnected = isConnected; // 添加!this.isManualDisconnect判断,主动断开后不触发重连 if (!this.reconnection$ && typeof(isConnected) === 'boolean' && !isConnected && !this.isManualDisconnect) { this.reconnect(); } });
重连方法中的循环条件
修改reconnect方法的takeWhile条件,加入主动断开标记,避免主动断开后仍继续重连:
private reconnect(): void { this.reconnection$ = interval(this.reconnectInterval) // 添加!this.isManualDisconnect,主动断开后停止重连尝试 .pipe(takeWhile((v, index) => index < this.reconnectAttempts && !this.websocket$ && !this.isManualDisconnect)); // 保存重连订阅,方便后续取消 this.reconnectionSubscription = this.reconnection$.subscribe( () => this.connect(), null, () => { this.reconnection$ = null; this.reconnectionSubscription = null; if (!this.websocket$) { this.wsMessages$.complete(); if (this.connection$) { this.connection$.complete(); } } }); }
连接方法中的前置判断
修改connect方法,主动断开状态下不再创建新连接:
private connect(): void { // 主动断开状态下,不再创建新连接 if (this.isManualDisconnect) return; this.websocket$ = new WebSocketSubject(this.config); this.websocket$.subscribe( (message) => this.wsMessages$.next(message), (error: Event) => { // 只有非主动断开时才触发重连 if (!this.websocket$ && !this.isManualDisconnect) { this.reconnect(); } }); }
4. 完善服务销毁逻辑
在ngOnDestroy中取消重连订阅,避免内存泄漏:
ngOnDestroy(): void { this.websocketSub.unsubscribe(); this.statusSub.unsubscribe(); if (this.reconnectionSubscription) { this.reconnectionSubscription.unsubscribe(); } }
内容的提问来源于stack exchange,提问作者Alex
相关产品推荐
相关产品推荐

