RxJS中如何实现Observable等待另一Observable值为true再执行逻辑
原有代码的问题
你当前写法的核心缺陷有两个:
B.filter(f => f)不会在取到第一个符合条件的值后终止,只要后续B再次变为true,之前所有A触发的内部订阅都会重复执行业务逻辑,不符合信号量单次放行的预期- 没有做流的终止处理,长期运行会出现内存泄漏、逻辑重复执行的问题
注意:不要使用withLatestFrom/combineLatest实现这个需求,这两个操作符会在A发值时如果B还没满足条件就直接丢弃当前A值,不符合“等待不丢值”的要求。
正确实现方式
核心逻辑是:A每发出一个值,就等待B发出第一个值为true的通知,拿到通知后立刻用当前A的值执行业务逻辑,同时结束当前内部的等待订阅。根据你需要的并发策略,选对应实现即可:
1. 并行放行
A连续发值时,只要B变为true,所有等待中的A值会同时执行,适合不需要排队的场景:
import { filter, first, map, mergeMap } from 'rxjs/operators'; A.pipe( mergeMap(aval => B.pipe( filter(bVal => bVal === true), first(), // 取到第一个true就终止当前等待,不会重复触发 map(() => aval) ) ) ).subscribe(aval => { // 执行业务逻辑 doSomething(aval) })
2. 排队放行(和C语言sem_wait逻辑完全一致)
A发出的值会按先后顺序排队,前一个值拿到信号量执行完之后,才会处理下一个等待的值,完全对应传统信号量的等待逻辑:
import { filter, first, map, concatMap } from 'rxjs/operators'; A.pipe( concatMap(aval => B.pipe( filter(bVal => bVal === true), first(), map(() => aval) ) ) ).subscribe(aval => { doSomething(aval) })
行为说明
- A发值时如果B当前值已经是
true:过滤逻辑会立刻放行,业务逻辑无延迟执行 - A发值时如果B当前值是
false:当前A值会进入等待状态,直到B下一次调用next(true)时才会触发执行 - 所有A发出的值都不会丢失,会按你选的策略等待信号量放行
- 每次等待拿到放行信号后就会终止当前内部订阅,不会因为B后续反复切换true/false重复执行业务
如果需要实现完整的信号量占用/释放逻辑,只需要在业务逻辑执行完成后手动调用
B.next(false)即可,对应C语言中信号量的sem_post操作。
内容的提问来源于stack exchange,提问作者Emme Developer
相关产品推荐
相关产品推荐

