Angular 15 WebSocket服务实现报错,寻求更优实现方案
问题解决与Angular 15中WebSocket服务的优化实现
一、当前代码报错修复
你的报错根源有两个:
- 使用RxJS内部类
AnonymousSubject:该类属于RxJS内部实现,不应直接从rxjs/internal/Subject导入,API不稳定且类型检查严格。 - 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(); } }
优化特性:
- 自动重连机制:通过
retryWhen和delayWhen实现断开后自动重试,提升服务稳定性。 - 连接状态通知:
isConnected$Subject向组件推送连接状态,便于UI展示。 - 类型安全:严格定义消息类型,避免JSON解析错误。
- 生命周期管理:实现
OnDestroy接口,销毁时关闭连接,防止内存泄漏。 - 清晰的职责划分:连接逻辑、消息发送、状态管理分离,代码更易维护。
组件使用示例
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
相关产品推荐
相关产品推荐

