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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 09:05:22