Angular NGRX Effect中使用withLatestFrom处理SSE流仅首次生效求助
问题原因与解决方案
核心问题排查
你的问题和NGRX是否支持RXJS操作符无关,完全是流逻辑编写错误导致SSE流被提前终止,具体问题如下:
- 第一个
map操作符没有返回值,导致SSE推送的EventDTO在流中丢失,后续withLatestFrom组合的第一个值为undefined。同时你使用类实例变量this.eventDto存储事件,在多值流场景下存在竞态覆盖风险。 - 第二个
map操作符存在分支遗漏:当事件类型、操作类型不匹配时没有返回有效的Action对象,NGRX Effect默认要求返回合法Action,收到undefined/非Action值会抛出错误。你配置的catchError捕获错误后返回EMPTY,会直接终止SSE对应的内部流,因此后续服务端推送不会再有任何响应。
修正后的实现
1. 调整Effect逻辑
不要使用实例变量存储事件,直接通过流传递SSE返回值,确保所有分支都有合法返回,避免流被意外终止:
loadEVents3$ = createEffect(() => this.actions$.pipe( ofType('[Websocketevent] Load Websocketevents'), mergeMap(() => this.webSocketEventsService.getAllNotify(`${environment.SERVER_API_URL}/services/my-api-events/api/event`).pipe( // 直接组合SSE事件与Store最新值,不需要中间变量存储 withLatestFrom( this.idAccountSelector$, this.idCompanySelector$, ), // 用mergeMap处理分支,不需要处理的场景直接返回EMPTY,避免发射无效值 mergeMap(([eventDto, idAccount, idCompany]) => { if (eventDto.eventType !== EventType.FOLDER_NOTE) { return EMPTY; } switch (eventDto.operationType) { case OperationType.CREATE: return of({ type: noteEventsActions.ActionNotesType.LOAD_WEB_SOCKET_EVENTS_NOTES, payload: eventDto }); case OperationType.UPDATE: return of({ type: noteEventsActions.ActionNotesType.UPDATE_WEB_SOCKET_EVENTS_NOTES, payload: eventDto }); case OperationType.DELETE: return of({ type: noteEventsActions.ActionNotesType.DELETE_WEB_SOCKET_EVENTS_NOTES, payload: eventDto }); default: return EMPTY; } }), catchError((err) => { console.error('SSE连接异常:', err); // 如果需要出错自动重连,可替换为 return throwError(() => err).pipe(retry(3)) return EMPTY; }) ) ) ) );
如果你希望保留map写法,也可以对不匹配的分支返回无操作的NOOP Action,再通过filter过滤掉即可。
2. 补充SSE连接销毁逻辑(可选,推荐)
你当前的SSE Observable没有实现取消逻辑,当订阅被终止(比如重复触发启动Action、页面销毁)时,EventSource不会自动关闭,会造成内存泄漏,修改getAllNotify方法补充销毁逻辑:
getAllNotify(url: string): Observable<EventDTO> { return new Observable(observer => { const eventSource = this.getEventSource(url); eventSource.onopen = function () { console.log('connection for all data is established'); }; eventSource.onmessage = event => { this.zone.run(() => { if (eventSource.readyState !== 0) { observer.next(JSON.parse(event.data)); } }); }; eventSource.onerror = error => { this.zone.run(() => { observer.error(error); }); }; // 订阅取消时自动关闭SSE连接 return () => eventSource.close(); }); }
额外优化建议
如果要避免重复启动多个SSE连接,可以将mergeMap替换为exhaustMap,前一个SSE连接未终止时,新的启动Action会被自动忽略。
内容的提问来源于stack exchange,提问作者Andy88
相关产品推荐
相关产品推荐

