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订阅连接。
具体实现步骤
- 在SSE对应的控制器中创建一个
Subject实例,用来存储和转发事件 - 单独写一个用
@OnEvent装饰的方法,监听stock.created事件,收到事件后调用Subject的next()方法推送数据 @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
相关产品推荐
相关产品推荐

