Angular+NgRx项目中用RxJS WaitUntil等待Firm$非空的实现方案
解决方案
根据你的场景,核心是在Effect中实现等待NgRx selector返回非空数组后再继续处理数据流,同时避免中断整个SSE流。以下是两种贴合不同需求的实现方案:
场景1:SSE已连接,每条事件需等待Firm就绪后处理
如果SSE服务已经启动,但需要确保每条通知事件都在Firm$返回非空数组后再处理(比如Firm数据可能在SSE运行期间才加载完成),可以用switchMap嵌套等待逻辑:
import { filter, first, switchMap, takeUntil } from 'rxjs/operators'; import { createEffect, ofType } from '@ngrx/effects'; notificationEffect$ = createEffect(() => this.actions$.pipe( ofType(NotificationActions.startSse), switchMap(() => // 启动SSE连接 this.sseService.connect().pipe( // 对每条SSE事件,先等待Firm$返回非空数组 switchMap(event => this.store.select(selectFirm).pipe( // 过滤非空数组(加Array.isArray避免null/undefined) filter(firms => Array.isArray(firms) && firms.length > 0), // 只取第一个满足条件的值,避免持续订阅Firm$ first(), // 拿到就绪信号后,把原事件传递下去 map(() => event) ) ), // 处理通知事件 map(event => NotificationActions.processNotification({ payload: event })), // 可选:监听停止动作,取消SSE订阅防止内存泄漏 takeUntil(this.actions$.pipe(ofType(NotificationActions.stopSse))) ) ) ) );
场景2:先等Firm就绪,再启动SSE连接
如果必须等Firm数据加载完成后才建立SSE连接(比如SSE依赖Firm的参数),可以把等待逻辑放在SSE连接之前:
import { filter, first, switchMap, takeUntil } from 'rxjs/operators'; notificationEffect$ = createEffect(() => this.actions$.pipe( ofType(NotificationActions.startSse), // 先等待Firm$返回非空数组 switchMap(() => this.store.select(selectFirm).pipe( filter(firms => Array.isArray(firms) && firms.length > 0), first() ) ), // Firm就绪后再启动SSE switchMap(() => this.sseService.connect().pipe( map(event => NotificationActions.processNotification({ payload: event })), takeUntil(this.actions$.pipe(ofType(NotificationActions.stopSse))) ) ) ) );
为什么之前的repeat/filter会出问题
- 直接在SSE流后加
filter:当Firm$为空时,所有SSE事件都会被过滤丢弃,即使后续Firm$变为非空,之前的事件也无法恢复。 - 使用
repeat:会让整个Effect流重新订阅,导致SSE服务重新连接,不符合"暂停流而非重启"的需求。
注意事项
- 用
first()确保拿到Firm就绪信号后立即取消订阅,避免不必要的内存占用。 - 加上
takeUntil监听停止动作,确保在关闭SSE时清理所有订阅,防止内存泄漏。
内容的提问来源于stack exchange,提问作者Andy88
相关产品推荐
相关产品推荐

