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

NestJS中如何结合EventEmitter与SSE推送后端事件至前端

问题:NestJS中通过SSE推送EventEmitter触发的后端事件

我需要将后端产生的事件单向推送到前端,由于是单向通信,不需要使用Socket,因此尝试用NestJS的SSE(Server Sent Events)实现需求。

按NestJS官方文档写的基础SSE代码可以正常运行:

@Sse('sse')
sse(): Observable<MessageEvent> {
  return interval(1000).pipe(map((_) => ({ data: { hello: 'world' } })));
}

但我需要把定时推送替换成推送后端实际产生的事件。目前已经实现了StocksService,创建Stock时会通过EventEmitter触发stock.created事件,且监听器能正常打印该事件:

StocksService代码

@Injectable()
export class StocksService {
  public stocks: Stock[] = [
    {
      id: 1,
      symbol: 'Stock #1',
      bid: 500,
      ask: 500,
    }
  ];

  constructor(private eventEmitter: EventEmitter2) {}

  create(createStockDto: CreateStockDto) {
    const stock = {
      id: this.stocks.length + 1,
      ...createStockDto,
    };
    this.stocks.push(stock);

    const stockCreatedEvent = new StockCreatedEvent();
    stockCreatedEvent.symbol = stock.symbol;
    stockCreatedEvent.ask = stock.ask;
    stockCreatedEvent.bid = stock.bid;

    this.eventEmitter.emit('stock.created', stockCreatedEvent);

    return stock;
  }
}

事件监听器代码

@Injectable()
export class StockCreatedListener {
  @OnEvent('stock.created')
  handleStockCreatedEvent(event: StockCreatedEvent) {
    console.log(event);
  }
}

但尝试把@Sse装饰器和@OnEvent事件监听结合时,访问http://localhost:3000/sse没有任何响应,不清楚该用Observable还是Subject来连接EventEmitter和SSE,求解决办法。

有问题的代码

@Sse('sse')
@OnEvent('stock.created')
sse(event: StockCreatedEvent): Observable<MessageEvent> {
  const obj = of(event);
  return obj.pipe(map((_) => ({ data: event})));
}

解决方案

核心是用Subject作为EventEmitter和SSE的桥梁,它既是Observable也是Observer,既能接收EventEmitter的事件,也能持续推送给前端的SSE订阅连接。

具体实现步骤

  1. 在SSE对应的控制器中创建一个Subject实例,用来存储和转发事件
  2. 单独写一个用@OnEvent装饰的方法,监听stock.created事件,收到事件后调用Subject的next()方法推送数据
  3. @Sse装饰的方法返回这个Subject,并转换成符合SSE要求的MessageEvent格式

完整控制器代码

import { Controller, Sse } from '@nestjs/common';
import { MessageEvent } from '@nestjs/common';
import { OnEvent } from '@nestjs/event-emitter';
import { Subject, Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import { StockCreatedEvent } from './events/stock-created.event';

@Controller()
export class SseController {
  // 创建Subject作为事件中转的桥梁
  private readonly stockEvents$ = new Subject<StockCreatedEvent>();

  @Sse('sse')
  sse(): Observable<MessageEvent> {
    // 将Subject转换为SSE要求的MessageEvent格式
    return this.stockEvents$.pipe(
      map(event => ({ data: event } as MessageEvent))
    );
  }

  @OnEvent('stock.created')
  handleStockCreatedEvent(event: StockCreatedEvent) {
    // 收到EventEmitter的事件后,推送给Subject
    this.stockEvents$.next(event);
  }
}

为什么之前的代码无效

之前把@Sse和@OnEvent放在同一个方法上是错误的:

  • @Sse修饰的方法需要返回一个持续存在的Observable来维持SSE连接,前端连接时这个方法会执行一次
  • @OnEvent修饰的方法是事件触发时才会执行的回调,两者生命周期不匹配
  • 之前用of(event)返回的Observable会在发送一次数据后立即完成,导致SSE连接直接关闭,所以前端看不到任何响应

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 14:35:12