如何使用RxJS Observables实现类广度优先搜索的异步队列逻辑
RxJS 类BFS动态队列实现方案
你可以使用expand+first+defaultIfEmpty的操作符组合实现需求,无需使用queueScheduler(queueScheduler仅用于控制同步任务调度顺序,不支持动态Observable队列的流控制、未执行任务销毁等能力)。
实现代码示例
import { of, Observable, EMPTY } from 'rxjs'; import { expand, first, defaultIfEmpty, delay } from 'rxjs/operators'; // 模拟你的内层业务Observable function innerTask(val: number): Observable<number> { // 模拟异步业务执行 return of(val).pipe(delay(100)); } // 构造外层Observable const bfsOuter$ = innerTask(1).pipe( // 动态追加内层任务到队列,完全匹配BFS遍历逻辑 expand(current => { // 业务逻辑:当前值小于3时追加2个子任务到队列 if (current < 3) { return of(innerTask(current * 2), innerTask(current * 2 + 1)); } // 返回EMPTY表示不再追加新任务 return EMPTY; }), // 命中目标立刻终止整个流,自动取消所有未执行任务 first( val => val === 5, // 命中规则,可替换为你的业务判断逻辑 val => val, // 命中后返回的值 null // 未命中时的默认返回值 ), // 兜底所有任务执行完无命中时返回null defaultIfEmpty(null) ); // 订阅使用 bfsOuter$.subscribe({ next: res => console.log('执行结果', res), complete: () => console.log('外层流结束') });
需求匹配说明
- 队列执行逻辑:
expand默认采用并发数为1的FIFO队列执行所有内层Observable,和BFS的遍历顺序完全一致,你可以在回调中根据当前内层执行结果动态追加任意数量的新内层Observable - 命中即终止:
first操作符匹配到目标值后会立即取消上游所有订阅,队列中未执行的内层Observable会被直接销毁,不会产生额外执行 - 无命中返回null:
first的第三个参数和defaultIfEmpty双重兜底,所有任务执行完未命中目标时自动返回null并结束外层流
内容的提问来源于stack exchange,提问作者argetlam5
相关产品推荐
相关产品推荐

