基于NestJS+React的SSE事件推送用户隔离问题求助
SSE事件推送范围问题解决(仅用SSE)
问题场景
- API A:基于NestJS开发,采用JWT认证策略
- API B:第三方不可控API,通过webhook向API A推送资源关联事件(如服务器宕机)
- 前端:React/NextJS开发,通过SSE连接API A接收事件
- 核心问题:API B推送某用户资源的事件时,所有在线用户都会收到通知,需实现仅资源所属用户收到对应事件
当前实现代码
NestJS Controller
@UseGuards(AuthGuard) @Sse('sse') sse(@Request() req) { return this.monitoringServerService.sendEvents(); }
MonitoringServerService
@Injectable() export class MonitoringService { private events = new Subject<any>(); addEvent(event) { this.events.next(event); } sendEvents() { return this.events.asObservable(); } }
解决方案
核心思路是按用户维护独立的SSE订阅流,替代全局Subject的广播模式。
1. 修改MonitoringService,维护用户订阅映射
@Injectable() export class MonitoringService { // 存储用户ID与专属事件流的映射关系 private userSubscriptions: Map<string, Subject<any>> = new Map(); // 获取或创建用户专属的事件流 getUserEventStream(userId: string): Observable<any> { if (!this.userSubscriptions.has(userId)) { this.userSubscriptions.set(userId, new Subject<any>()); } return this.userSubscriptions.get(userId).asObservable(); } // 向指定用户推送事件 addEventToUser(userId: string, event: any) { const userStream = this.userSubscriptions.get(userId); if (userStream) { userStream.next(event); } } // 用户断开连接时清理订阅,避免内存泄漏 removeUserSubscription(userId: string) { const userStream = this.userSubscriptions.get(userId); if (userStream) { userStream.complete(); this.userSubscriptions.delete(userId); } } }
2. 修改SSE Controller,绑定用户专属流
添加断开连接的处理逻辑:
@Controller() export class SseController { constructor(private readonly monitoringService: MonitoringService) {} @UseGuards(AuthGuard) @Sse('sse') sse(@Request() req): Observable<any> { // 从JWT解析当前用户ID(根据你的AuthGuard实现调整,比如req.user.id) const userId = req.user.id; return this.monitoringService.getUserEventStream(userId); } @OnDisconnect() handleDisconnect(@Request() req) { const userId = req.user.id; this.monitoringService.removeUserSubscription(userId); } }
3. 处理API B的Webhook请求
在接收webhook的控制器中,解析事件所属用户ID,定向推送:
@Controller('webhook') export class WebhookController { constructor(private readonly monitoringService: MonitoringService) {} @Post('event') handleWebhook(@Body() eventData: any) { // 从API B的事件数据中提取资源所属的用户ID(根据API B返回结构调整) const userId = eventData.resource.ownerId; // 向指定用户推送事件 this.monitoringService.addEventToUser(userId, { type: 'resource-event', data: eventData }); return { status: 'ok' }; } }
关键说明
- 每个用户连接SSE时会创建独立的Subject,仅该用户能收到对应流的事件
- 用户断开连接时必须清理订阅,防止内存泄漏
- Webhook处理环节需确保能从API B的事件数据中准确识别所属用户ID
内容的提问来源于stack exchange,提问作者yuridamata
相关产品推荐
相关产品推荐

