在Nest网关中如何停止旧Observable并在事件后创建新实例?
解决方案:管理WebSocket连接的订阅生命周期
问题核心是每次调用listenAs都会生成新的Observable订阅,但旧订阅未被取消,导致多个订阅同时生效,客户端会收到多组过滤后的消息。要解决这个问题,你需要在网关层面为每个WebSocket连接维护订阅实例,确保切换监听用户时先终止旧订阅,再启动新订阅。
步骤1:在网关中维护订阅映射
在WebSocket网关内添加一个Map,存储每个客户端的当前订阅(键为WebSocket实例,值为RxJS的Subscription对象),同时处理连接断开时的订阅清理:
import { WebSocketGateway, WebSocketServer, SubscribeMessage, OnGatewayConnection, OnGatewayDisconnect, } from '@nestjs/websockets'; import { Server, WebSocket } from 'ws'; import { Subscription } from 'rxjs'; import { MessagesService } from './messages.service'; import { ListenAsDto } from './dto/listen-as.dto'; @WebSocketGateway({ transport: 'ws' }) export class MessagesGateway implements OnGatewayConnection, OnGatewayDisconnect { @WebSocketServer() server: Server; // 存储每个客户端的活跃订阅 private clientSubscriptions = new Map<WebSocket, Subscription>(); constructor(private readonly messagesService: MessagesService) {} // 客户端断开时清理订阅 handleDisconnect(client: WebSocket) { const subscription = this.clientSubscriptions.get(client); if (subscription) { subscription.unsubscribe(); this.clientSubscriptions.delete(client); } } @SubscribeMessage('listen-as') handleListenAs(client: WebSocket, payload: ListenAsDto) { // 取消当前客户端的旧订阅 const existingSubscription = this.clientSubscriptions.get(client); if (existingSubscription) { existingSubscription.unsubscribe(); } // 创建新订阅并向客户端推送消息 const newSubscription = this.messagesService.listenAs(payload).subscribe({ next: (response) => { client.send(JSON.stringify(response)); }, error: (err) => { client.send(JSON.stringify({ event: 'error', data: err.message })); }, }); // 更新映射中的当前订阅 this.clientSubscriptions.set(client, newSubscription); } }
步骤2:优化消息数据流(可选但推荐)
如果你的messageObservable是普通Subject,每次订阅会触发独立数据流,建议改用share()操作符让多个订阅共享同一数据流,避免重复处理消息:
import { Injectable } from '@nestjs/common'; import { Subject, Observable } from 'rxjs'; import { share } from 'rxjs/operators'; import { WsResponse } from '@nestjs/websockets'; import { ListenAsDto } from './dto/listen-as.dto'; import { PersonalisedMessage } from './interfaces/personalised-message.interface'; @Injectable() export class MessagesService { private messageSubject = new Subject<{ to: string } & PersonalisedMessage>(); // 共享数据流,减少重复处理 messageObservable = this.messageSubject.asObservable().pipe(share()); // 内部调用此方法推送消息 sendMessage(to: string, message: PersonalisedMessage) { this.messageSubject.next({ ...message, to }); } listenAs(listener: ListenAsDto): Observable<WsResponse<PersonalisedMessage>> { return this.messageObservable.pipe( filter((item) => item.to === listener.name), map((item) => { const { to, ...rest } = item; return { event: 'INBOX_MESSAGE_NAME', data: rest }; }), ); } }
关键说明
- 单订阅约束:每次处理
listen-as事件时,先终止客户端的旧订阅,确保同一时间每个客户端只有一个活跃订阅,避免接收多用户消息。 - 内存泄漏防护:实现
OnGatewayDisconnect接口,在客户端断开时清理订阅,防止无用订阅占用资源。 - 性能优化:通过
share()让多个订阅共享数据流,避免重复处理相同消息,提升服务器性能。
内容的提问来源于stack exchange,提问作者RichGihratik
相关产品推荐
相关产品推荐

