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

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;
  }
}

关键注意事项

  1. 用户认证:必须确保SSE请求能获取到用户的身份/角色信息(比如通过JWT token放在请求头,或者会话机制),这是定向推送的前提。
  2. 内存泄漏:方案二中一定要监听连接关闭事件,调用unregisterUserStream清理Subject,避免内存泄漏。
  3. SSE格式要求:返回的Observable必须是{ data: any }结构,这是SSE协议的要求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 16:45:01