如何用RxJS现有运算符实现递归switchMap链式获取nextState$状态流
可行方案:使用RxJS原生
expand运算符实现 完全可以通过RxJS内置运算符实现,不需要自定义Observable,核心使用专门用于递归处理流的expand运算符即可同时满足你的两个需求。
核心实现代码
import { of, expand, EMPTY } from 'rxjs'; // 初始状态 const initialState = { nextState$: createObservableWithNextState$() }; const states$ = of(initialState).pipe( // 递归处理每个状态的nextState$ expand(state => { // 若当前状态不存在nextState$则返回空Observable终止递归 return state.nextState$ ?? EMPTY; }) );
效果说明
上述实现完美匹配你的需求:
- 无需重复编写
switchMap:expand会自动递归处理每一个发射出来的状态对象,不需要手动堆叠运算符 - 能收到全量状态:源Observable发射的初始状态会优先传递给下游订阅者,后续每一轮
nextState$推送的状态也会依次输出,不会丢失任何节点
可运行测试示例
import { of, expand, delay } from 'rxjs'; // 模拟状态生成逻辑,最多生成5个状态 let stateIndex = 0; function createObservableWithNextState$() { stateIndex++; if (stateIndex >= 5) return null; // 模拟异步延迟推送下一个状态 return of({ id: stateIndex, nextState$: createObservableWithNextState$() }).pipe(delay(1000)); } // 初始化流 const states$ = of({ id: 0, nextState$: createObservableWithNextState$() }).pipe( expand(state => state.nextState$ ?? []) ); // 订阅输出 states$.subscribe(state => console.log('收到状态ID:', state.id));
运行后输出顺序为:
收到状态ID:0 // 初始状态立刻输出 收到状态ID:1 // 1秒后输出 收到状态ID:2 // 2秒后输出 收到状态ID:3 // 3秒后输出 收到状态ID:4 // 4秒后输出
原理解释
expand是RxJS专门用于递归展开高阶Observable的运算符,工作逻辑如下:
- 上游每推送一个值,会先将该值直接发送给下游订阅者
- 调用你传入的回调函数,将当前值作为入参,拿到返回的Observable
- 自动订阅这个Observable,将它推送的所有值作为新的上游输入,重复执行上述流程,直到回调返回的Observable触发complete事件,整个流自动终止
如果需要和你原有示例的switchMap行为完全对齐(即新的状态生成后自动取消上一个未完成的nextState$订阅),可以给expand添加第三个参数concurrency设置为1,常规场景下上述默认实现已经满足需求。
内容的提问来源于stack exchange,提问作者Charles HETIER
相关产品推荐
相关产品推荐

