如何在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
相关产品推荐
相关产品推荐

