NestJS中@Sse实现定向推送:如何按用户角色发送专属通知?
解决NestJS中SSE按用户角色定向推送专属通知的问题
核心问题分析
你当前的代码问题在于每个SSE连接都会创建独立的Observable实例,如果在Observable内部用setInterval定时推送,就会导致每个用户连接都启动一个定时器,且消息无差别广播。要实现定向推送,核心是把「消息发送逻辑」和「客户端订阅逻辑」分离,通过用户身份/角色做过滤。
方案一:用全局Subject+角色过滤(基于Observable实现)
这是最贴合NestJS SSE设计的方案,通过全局消息流结合RxJS的filter操作符,让客户端只接收符合自身角色的消息。
1. 创建SSE管理服务
// sse.service.ts import { Injectable } from '@nestjs/common'; import { Subject, filter, map } from 'rxjs'; // 定义定向消息结构:包含目标角色和内容 interface TargetedMessage { targetRoles: string[]; content: any; } @Injectable() export class SseService { // 全局消息主题,所有消息都通过它分发 private messageSubject = new Subject<TargetedMessage>(); // 对外暴露的发送方法:指定目标角色和消息内容 sendToRoles(targetRoles: string[], content: any) { this.messageSubject.next({ targetRoles, content }); } // 给单个用户返回专属的Observable流 getUserStream(userRoles: string[]) { return this.messageSubject.pipe( // 只保留用户角色包含在目标角色里的消息 filter(msg => msg.targetRoles.some(role => userRoles.includes(role))), // 转换为SSE要求的格式(必须包含data字段) map(msg => ({ data: msg.content })) ); } }
2. 实现SSE控制器
// sse.controller.ts import { Controller, Sse, Req } from '@nestjs/common'; import { Request } from 'express'; import { Observable } from 'rxjs'; import { SseService } from './sse.service'; @Controller('sse') export class SseController { constructor(private readonly sseService: SseService) {} @Sse('stream') getStream(@Req() req: Request): Observable<{ data: any }> { // 假设你已经通过JWT/会话完成用户认证,req.user包含用户角色信息 const userRoles = req.user.roles; // 返回该用户专属的过滤后的流 return this.sseService.getUserStream(userRoles); } }
3. 业务逻辑中发送定向消息
在任意服务中注入SseService,即可给指定角色发消息:
// 示例:订单服务 import { Injectable } from '@nestjs/common'; import { SseService } from './sse.service'; @Injectable() export class OrderService { constructor(private readonly sseService: SseService) {} async createOrder() { // 业务逻辑... // 给管理员发送订单创建通知 this.sseService.sendToRoles(['admin'], { type: 'order_created', content: '新订单已创建,请及时处理', time: new Date().toISOString() }); // 给下单用户发送确认通知 this.sseService.sendToRoles(['user'], { type: 'order_confirm', content: '你的订单已提交成功', time: new Date().toISOString() }); } }
方案二:单用户独立流(更细粒度控制)
如果需要给单个用户推送专属消息,可以用Map存储每个用户的独立Subject,实现精准推送。
1. 修改SSE服务
// sse.service.ts import { Injectable } from '@nestjs/common'; import { Subject, map } from 'rxjs'; // 存储单个用户的流信息 interface UserStream { subject: Subject<any>; roles: string[]; } @Injectable() export class SseService { private userStreams = new Map<string, UserStream>(); // key: 用户ID // 用户连接时注册专属流 registerUserStream(userId: string, roles: string[]) { const subject = new Subject<any>(); this.userStreams.set(userId, { subject, roles }); return subject.pipe(map(content => ({ data: content }))); } // 用户断开连接时清理流 unregisterUserStream(userId: string) { const stream = this.userStreams.get(userId); if (stream) { stream.subject.complete(); this.userStreams.delete(userId); } } // 给单个用户发消息 sendToUser(userId: string, content: any) { const stream = this.userStreams.get(userId); stream?.subject.next(content); } // 给指定角色发消息 sendToRoles(targetRoles: string[], content: any) { this.userStreams.forEach((stream) => { if (stream.roles.some(role => targetRoles.includes(role))) { stream.subject.next(content); } }); } }
2. 控制器中处理连接生命周期
// sse.controller.ts import { Controller, Sse, Req } from '@nestjs/common'; import { Request } from 'express'; import { Observable } from 'rxjs'; import { SseService } from './sse.service'; @Controller('sse') export class SseController { constructor(private readonly sseService: SseService) {} @Sse('stream') getStream(@Req() req: Request): Observable<{ data: any }> { const userId = req.user.id; const userRoles = req.user.roles; // 注册用户流 const stream = this.sseService.registerUserStream(userId, userRoles); // 监听连接关闭事件,清理资源 req.on('close', () => { this.sseService.unregisterUserStream(userId); }); return stream; } }
关键注意事项
- 用户认证:必须确保SSE请求能获取到用户的身份/角色信息(比如通过JWT token放在请求头,或者会话机制),这是定向推送的前提。
- 内存泄漏:方案二中一定要监听连接关闭事件,调用
unregisterUserStream清理Subject,避免内存泄漏。 - SSE格式要求:返回的Observable必须是
{ data: any }结构,这是SSE协议的要求。
内容的提问来源于stack exchange,提问作者Bennison J
相关产品推荐
相关产品推荐

