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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 18:29:55