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

Angular中如何结合RxJS实现WebSocket的条件化禁用重试

解决方案

问题根源在于你直接将retry操作符附加到了WebSocketSubject上,并且将管道后的Observable断言为WebSocketSubject。这会导致两个核心问题:

  1. 管道后的Observable并非真正的WebSocketSubject,调用complete()无法正确关闭底层连接
  2. 主动关闭时,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!`
                    );
                }
            }
        };
    }
}

关键修改点说明

  1. 分离Subject与Observable:

    • 原始WebSocketSubject用于直接控制连接的打开/关闭
    • 带retryWhen逻辑的Observable用于订阅消息流,避免重试逻辑干扰连接控制
  2. 主动关闭标记:

    • 调用close()时标记对应Socket为主动关闭,retryWhen检测到该标记后终止重试流程
    • 重新调用connect()时重置标记,恢复自动重试能力
  3. 灵活的重试控制:

    • 使用retryWhen替代retry,可以根据重试次数、关闭类型等条件决定是否继续重试
    • 增加了重试日志,便于排查连接问题

内容的提问来源于stack exchange,提问作者RAHUL KUNDU

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 04:59:51