Angular中如何结合RxJS实现WebSocket的条件化禁用重试
解决方案
问题根源在于你直接将retry操作符附加到了WebSocketSubject上,并且将管道后的Observable断言为WebSocketSubject。这会导致两个核心问题:
- 管道后的Observable并非真正的
WebSocketSubject,调用complete()无法正确关闭底层连接 - 主动关闭时,
retry会将关闭事件视为错误并触发重连
核心修改思路
- 为每个Socket连接维护主动关闭标记,区分主动关闭和意外断开场景
- 分离原始
WebSocketSubject和带重试逻辑的Observable,确保complete()直接作用于原始连接 - 使用
retryWhen替代retry,实现更灵活的重试控制逻辑
修改后的完整代码
import { NextObserver, Observable, retryWhen, mergeMap, throwError, timer } from 'rxjs'; import { Injectable } from '@angular/core'; import { appSettings } from '@app/configs'; import { environment } from '@env/environment'; import { ISocketResult } from '@shared/models'; import { LoggerService } from './logger.service'; import { webSocket, WebSocketSubject, WebSocketSubjectConfig } from 'rxjs/webSocket'; @Injectable() export class SocketService { constructor(private _logger: LoggerService) {} private isAdminSocketConnected = false; private isFrontendSocketConnected = false; private retryDelay = appSettings.socketRetryDelayMS; private WS_ADMIN_ENDPOINT = environment.adminSocketHost; private maxRetryAttempts = appSettings.socketRetryAttempts; private WS_FRONTEND_ENDPOINT = environment.frontendSocketHost; // 存储原始WebSocketSubject,用于直接控制连接关闭 private _adminSocketSubject!: WebSocketSubject<ISocketResult | null>; private _frontendSocketSubject!: WebSocketSubject<ISocketResult | null>; // 存储带重试逻辑的Observable,用于订阅消息 private _adminSocket$!: Observable<ISocketResult | null>; private _frontendSocket$!: Observable<ISocketResult | null>; // 标记是否主动关闭对应Socket private _isSocketManuallyClosed: Record<string, boolean> = {}; public connect(): void { // 连接前端Socket if (!this._frontendSocketSubject || this._frontendSocketSubject.closed) { this._isSocketManuallyClosed[this.WS_FRONTEND_ENDPOINT] = false; const { subject, observable } = this.createSocketInstance(this.WS_FRONTEND_ENDPOINT); this._frontendSocketSubject = subject; this._frontendSocket$ = observable; this.subscribeToSocket(this._frontendSocket$, 'frontend'); } // 连接管理员Socket if (!this._adminSocketSubject || this._adminSocketSubject.closed) { this._isSocketManuallyClosed[this.WS_ADMIN_ENDPOINT] = false; const { subject, observable } = this.createSocketInstance(this.WS_ADMIN_ENDPOINT); this._adminSocketSubject = subject; this._adminSocket$ = observable; this.subscribeToSocket(this._adminSocket$, 'admin'); } } public close(): void { // 关闭管理员Socket this._isSocketManuallyClosed[this.WS_ADMIN_ENDPOINT] = true; this._adminSocketSubject?.complete(); this._adminSocketSubject = undefined!; // 关闭前端Socket this._isSocketManuallyClosed[this.WS_FRONTEND_ENDPOINT] = true; this._frontendSocketSubject?.complete(); this._frontendSocketSubject = undefined!; } private createSocketInstance( endpoint: string ): { subject: WebSocketSubject<ISocketResult | null>, observable: Observable<ISocketResult | null> } { const socketConfig: WebSocketSubjectConfig<ISocketResult | null> = { url: endpoint, openObserver: this.createOpenObserver(endpoint), closeObserver: this.createCloseObserver(endpoint), deserializer: (event: MessageEvent) => { try { const data = JSON.parse(event.data) as ISocketResult; return data; } catch (error) { console.error('Error during deserialization:', error); return null; } } }; const subject = webSocket(socketConfig); // 使用retryWhen实现带条件的重试 const observable = subject.pipe( retryWhen(errors => errors.pipe( mergeMap((error, retryIndex) => { // 主动关闭时终止重试 if (this._isSocketManuallyClosed[endpoint]) { return throwError(() => new Error(`Socket ${endpoint} manually closed`)); } // 达到最大重试次数时终止 if (retryIndex >= this.maxRetryAttempts) { this._logger.error(`Max retry attempts (${this.maxRetryAttempts}) reached for ${endpoint}`); return throwError(() => new Error(`Max retry attempts reached`)); } // 延迟后重试 this._logger.warn(`Retrying ${endpoint} connection (${retryIndex + 1}/${this.maxRetryAttempts})...`); return timer(this.retryDelay); }) ) ) ); return { subject, observable }; } private subscribeToSocket(socket$: Observable<ISocketResult | null>, type: string): void { socket$.subscribe({ next: (data) => { // 处理收到的消息,根据业务需求添加逻辑 if (data) { this._logger.info(`Received ${type} socket message:`, data); } }, error: (err) => { this._logger.error(`${type} socket error:`, err); }, complete: () => { this._logger.info(`${type} socket connection completed`); } }); } private createOpenObserver( endpoint: string ): NextObserver<Event | undefined> { return { next: (event) => { switch (endpoint) { case this.WS_ADMIN_ENDPOINT: if (!this.isAdminSocketConnected) this.isAdminSocketConnected = true; break; case this.WS_FRONTEND_ENDPOINT: if (!this.isFrontendSocketConnected) this.isFrontendSocketConnected = true; break; default: this._logger.error( `Unexpected WebSocket server at ${endpoint}. Unable to determine the source.` ); return; } if (this.isAdminSocketConnected && this.isFrontendSocketConnected) { this._logger.info( 'Frontend and Admin sockets are connected to the WebSocket server!' ); } } }; } private createCloseObserver( endpoint: string ): NextObserver<Event | undefined> { return { next: (event) => { switch (endpoint) { case this.WS_FRONTEND_ENDPOINT: if (this.isFrontendSocketConnected) this.isFrontendSocketConnected = false; break; case this.WS_ADMIN_ENDPOINT: if (this.isAdminSocketConnected) this.isAdminSocketConnected = false; break; default: this._logger.warn( `Unexpected WebSocket server at ${endpoint}.` ); return; } if (!this.isAdminSocketConnected && !this.isFrontendSocketConnected) { this._logger.info( `Frontend and Admin sockets are disconnected from the WebSocket server!` ); } } }; } }
关键修改点说明
分离Subject与Observable:
- 原始
WebSocketSubject用于直接控制连接的打开/关闭 - 带
retryWhen逻辑的Observable用于订阅消息流,避免重试逻辑干扰连接控制
- 原始
主动关闭标记:
- 调用
close()时标记对应Socket为主动关闭,retryWhen检测到该标记后终止重试流程 - 重新调用
connect()时重置标记,恢复自动重试能力
- 调用
灵活的重试控制:
- 使用
retryWhen替代retry,可以根据重试次数、关闭类型等条件决定是否继续重试 - 增加了重试日志,便于排查连接问题
- 使用
内容的提问来源于stack exchange,提问作者RAHUL KUNDU
相关产品推荐
相关产品推荐

