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

如何在NestJS中实现模块增删改时触发Server Sent Event?

在NestJS中实现模块增删改的SSE推送方案

问题描述

我想在模块创建、更新或删除时发送SSE事件,查了NestJS的SSE文档没找到有效帮助。下面是模块增删改的方法代码,我已经创建了控制器,这是第一次在NestJS里用SSE,想知道调用这些方法时实现SSE的最佳方案。

addNewModule(data: ApiGateway.Event.AddModule.IPayload) {
    this.logger.debug('Handle event on add new module!');
    this.logger.debug(`Previous modules state: ${this.modules}`);
    this.modules.push(data.moduleInfo);
    this.logger.debug(`New modules state: ${this.modules}`);
  }

  removeModule(data: ApiGateway.Event.RemoveModule.IPayload) {
    this.logger.debug('Handle event on add new module!');
    this.logger.debug(`Previous modules state: ${this.modules}`);
    this.modules = this.modules.filter(module => {
      return module.name !== data.moduleName;
    });
    this.logger.debug(`New modules state: ${this.modules}`);
  }

  updateModuleSettings(data: ApiGateway.Event.UpdateModuleSettings.IPayload) {
    this.logger.debug('Handle event on add new module!');
    this.logger.debug(`Previous modules state: ${this.modules}`);
    this.modules.find((module, index) => {
      if (module.name === data.moduleInfo.name) {
        this.modules[index] = data.moduleInfo;
        return true;
      }
    });
    this.logger.debug(`New modules state: ${this.modules}`);
  }

解决方案

核心思路是用RxJS的Subject做事件广播器,在增删改操作完成后触发事件,再通过控制器的SSE端点推送给客户端。

步骤1:改造服务层,添加事件发布逻辑

在你的模块服务里引入Subject,定义事件类型,在每个增删改方法末尾发布对应事件:

import { Injectable, Logger } from '@nestjs/common';
import { Subject } from 'rxjs';
import { Observable } from 'rxjs';

// 定义模块变更事件的统一格式
type ModuleChangeEvent = {
  type: 'add' | 'remove' | 'update';
  data: any; // 可替换成你实际的Payload类型
};

@Injectable()
export class YourModuleService {
  private readonly logger = new Logger(YourModuleService.name);
  private modules: any[] = [];
  // 创建Subject作为事件流载体
  private moduleChangeSubject = new Subject<ModuleChangeEvent>();

  // 对外暴露Observable,防止外部直接修改事件流
  getModuleChangeEvents(): Observable<ModuleChangeEvent> {
    return this.moduleChangeSubject.asObservable();
  }

  addNewModule(data: ApiGateway.Event.AddModule.IPayload) {
    this.logger.debug('处理新增模块事件!');
    this.logger.debug(`变更前模块状态: ${JSON.stringify(this.modules)}`);
    this.modules.push(data.moduleInfo);
    this.logger.debug(`变更后模块状态: ${JSON.stringify(this.modules)}`);
    // 发布新增事件
    this.moduleChangeSubject.next({
      type: 'add',
      data: data.moduleInfo,
    });
  }

  removeModule(data: ApiGateway.Event.RemoveModule.IPayload) {
    this.logger.debug('处理删除模块事件!');
    this.logger.debug(`变更前模块状态: ${JSON.stringify(this.modules)}`);
    this.modules = this.modules.filter(module => module.name !== data.moduleName);
    this.logger.debug(`变更后模块状态: ${JSON.stringify(this.modules)}`);
    // 发布删除事件
    this.moduleChangeSubject.next({
      type: 'remove',
      data: data.moduleName,
    });
  }

  updateModuleSettings(data: ApiGateway.Event.UpdateModuleSettings.IPayload) {
    this.logger.debug('处理更新模块配置事件!');
    this.logger.debug(`变更前模块状态: ${JSON.stringify(this.modules)}`);
    const targetIndex = this.modules.findIndex(module => module.name === data.moduleInfo.name);
    if (targetIndex !== -1) {
      this.modules[targetIndex] = data.moduleInfo;
    }
    this.logger.debug(`变更后模块状态: ${JSON.stringify(this.modules)}`);
    // 发布更新事件
    this.moduleChangeSubject.next({
      type: 'update',
      data: data.moduleInfo,
    });
  }
}

步骤2:创建SSE控制器端点

在控制器里注入服务,用@SSE装饰器创建事件订阅路由,把事件流转换成SSE格式返回:

import { Controller, SSE } from '@nestjs/common';
import { YourModuleService } from './your-module.service';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';

@Controller('modules')
export class ModulesController {
  constructor(private readonly moduleService: YourModuleService) {}

  @SSE('change-events')
  getModuleChangeEvents(): Observable<MessageEvent> {
    return this.moduleService.getModuleChangeEvents().pipe(
      map(event => ({
        type: `module.${event.type}`, // 自定义事件类型,方便客户端区分
        data: event.data,
      })),
    );
  }
}

步骤3:客户端订阅SSE

用浏览器原生的EventSource或者第三方库订阅端点,监听对应事件:

const eventSource = new EventSource('/modules/change-events');

// 监听新增模块事件
eventSource.addEventListener('module.add', (e) => {
  const newModule = JSON.parse(e.data);
  console.log('新增模块:', newModule);
});

// 监听删除模块事件
eventSource.addEventListener('module.remove', (e) => {
  const deletedModuleName = e.data;
  console.log('删除模块:', deletedModuleName);
});

// 监听更新模块事件
eventSource.addEventListener('module.update', (e) => {
  const updatedModule = JSON.parse(e.data);
  console.log('更新模块:', updatedModule);
});

关键注意事项

  • 用Subject而非BehaviorSubject:不需要给新订阅者发送历史事件,只推送实时变更。
  • 暴露Observable而非直接暴露Subject:避免外部代码随意触发事件,保证事件流的可控性。
  • NestJS的@SSE装饰器会自动设置text/event-stream响应头,无需手动配置。
  • 如需重连支持:客户端可添加断开重连逻辑,或在服务层定期发送心跳事件维持连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 19:10:34