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

Angular 15 WebSocket服务实现报错,寻求更优实现方案

问题解决与Angular 15中WebSocket服务的优化实现

一、当前代码报错修复

你的报错根源有两个:

  1. 使用RxJS内部类AnonymousSubject:该类属于RxJS内部实现,不应直接从rxjs/internal/Subject导入,API不稳定且类型检查严格。
  2. Observer未实现完整接口:你定义的observer对象中error和complete为null,不符合Observer接口要求(必须包含这两个方法)。

修复后的代码(保留原有逻辑)

// src\app\services\websocket.service.ts
import { Injectable } from "@angular/core";
import { Observable, Observer, Subject } from 'rxjs';
import { map } from 'rxjs/operators';

const CHAT_URL = "ws://localhost:5000";

export interface Message {
  source: string;
  content: string;
}

@Injectable({ providedIn: 'root' }) // 启用树摇优化,符合Angular最佳实践
export class WebsocketService {
  private subject: Subject<MessageEvent> | undefined;
  public messages: Subject<Message>;

  constructor() {
    this.messages = this.connect(CHAT_URL).pipe(
      map((response: MessageEvent): Message => {
        console.log(response.data);
        return JSON.parse(response.data);
      })
    ) as Subject<Message>;
  }

  public connect(url: string): Subject<MessageEvent> {
    if (!this.subject) {
      this.subject = this.create(url);
      console.log("Successfully connected: " + url);
    }
    return this.subject;
  }

  private create(url: string): Subject<MessageEvent> {
    const ws = new WebSocket(url);
    const observable = new Observable<MessageEvent>(obs => {
      ws.onmessage = obs.next.bind(obs);
      ws.onerror = obs.error.bind(obs);
      ws.onclose = obs.complete.bind(obs);
      return () => ws.close();
    });

    const observer: Observer<unknown> = {
      next: (data: unknown) => {
        console.log('Message sent to websocket: ', data);
        if (ws.readyState === WebSocket.OPEN) {
          ws.send(JSON.stringify(data));
        }
      },
      error: (err: unknown) => {
        console.error('WebSocket send error:', err);
        ws.readyState === WebSocket.OPEN && ws.close();
      },
      complete: () => {
        console.log('WebSocket observer completed');
        ws.readyState === WebSocket.OPEN && ws.close();
      }
    };

    return Subject.create(observer, observable);
  }
}

关键修复点:

  • 移除AnonymousSubject,改用RxJS公共APISubject.create()创建Subject,避免依赖内部类。
  • 补全observer的error和complete方法,满足接口规范。
  • 为@Injectable添加providedIn: 'root',实现服务的树摇优化。

二、Angular 15中更优的WebSocket服务实现

以下是更健壮、符合现代Angular规范的实现,包含自动重连、状态管理、类型安全等核心功能:

// src\app\services\websocket.service.ts
import { Injectable, OnDestroy } from "@angular/core";
import { Observable, Subject, timer } from 'rxjs';
import { retryWhen, delayWhen, map, tap, finalize } from 'rxjs/operators';

const CHAT_URL = "ws://localhost:5000";
const RECONNECT_DELAY = 3000; // 重连间隔(毫秒)

export interface Message {
  source: string;
  content: string;
}

@Injectable({ providedIn: 'root' })
export class WebsocketService implements OnDestroy {
  private wsSubject$: Subject<MessageEvent> | undefined;
  public messages$: Observable<Message>;
  public isConnected$ = new Subject<boolean>();

  constructor() {
    this.messages$ = this.connect().pipe(
      map(event => JSON.parse(event.data) as Message),
      tap(() => this.isConnected$.next(true)),
      // 断开后自动重连
      retryWhen(errors => 
        errors.pipe(
          tap(() => this.isConnected$.next(false)),
          delayWhen(() => timer(RECONNECT_DELAY))
        )
      ),
      finalize(() => this.isConnected$.next(false))
    );
  }

  private connect(): Observable<MessageEvent> {
    return new Observable(obs => {
      const ws = new WebSocket(CHAT_URL);

      ws.onopen = () => {
        console.log('WebSocket connected');
        this.isConnected$.next(true);
      };

      ws.onmessage = obs.next.bind(obs);
      ws.onerror = err => {
        console.error('WebSocket error:', err);
        obs.error(err);
      };
      ws.onclose = () => {
        console.log('WebSocket disconnected');
        obs.complete();
      };

      return () => {
        ws.readyState === WebSocket.OPEN && ws.close();
      };
    });
  }

  sendMessage(message: Message): void {
    if (this.wsSubject$) {
      this.wsSubject$.next(JSON.stringify(message));
    } else {
      console.warn('Cannot send message: WebSocket not connected');
    }
  }

  ngOnDestroy(): void {
    this.wsSubject$?.complete();
  }
}

优化特性:

  1. 自动重连机制:通过retryWhen和delayWhen实现断开后自动重试,提升服务稳定性。
  2. 连接状态通知:isConnected$ Subject向组件推送连接状态,便于UI展示。
  3. 类型安全:严格定义消息类型,避免JSON解析错误。
  4. 生命周期管理:实现OnDestroy接口,销毁时关闭连接,防止内存泄漏。
  5. 清晰的职责划分:连接逻辑、消息发送、状态管理分离,代码更易维护。

组件使用示例

import { Component, OnInit, OnDestroy } from '@angular/core';
import { WebsocketService, Message } from './services/websocket.service';
import { Subscription } from 'rxjs';

@Component({
  selector: 'app-chat',
  template: `
    <div class="status" [class.connected]="isConnected">
      {{isConnected ? '已连接' : '未连接,正在重试...'}}
    </div>
    <div class="messages">
      <div *ngFor="let msg of messages" class="message">
        <span class="source">{{msg.source}}:</span> {{msg.content}}
      </div>
    </div>
    <input [(ngModel)]="newMessage" (keyup.enter)="sendMessage()" placeholder="输入消息...">
  `,
  styles: [`
    .status.connected { color: green; }
    .messages { margin: 1rem 0; }
    .message { margin: 0.5rem 0; }
    .source { font-weight: bold; }
  `]
})
export class ChatComponent implements OnInit, OnDestroy {
  messages: Message[] = [];
  newMessage = '';
  isConnected = false;
  private subs = new Subscription();

  constructor(private wsService: WebsocketService) {}

  ngOnInit(): void {
    this.subs.add(
      this.wsService.messages$.subscribe(msg => {
        this.messages.push(msg);
      })
    );

    this.subs.add(
      this.wsService.isConnected$.subscribe(status => {
        this.isConnected = status;
      })
    );
  }

  sendMessage(): void {
    if (this.newMessage.trim()) {
      this.wsService.sendMessage({
        source: 'user',
        content: this.newMessage.trim()
      });
      this.newMessage = '';
    }
  }

  ngOnDestroy(): void {
    this.subs.unsubscribe();
  }
}

内容的提问来源于stack exchange,提问作者shujaat siddiqui

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 14:53:19