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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 01:22:29