RxJS可观察对象优先级队列实现:低优先级任务延迟发射方案
RxJS 实现优先级队列:低优先级Observable等待高优先级凑组或超时后发射
我刚好搞定了这个RxJS里的优先级队列问题,核心就是用timeoutWith结合嵌套的zip操作符,完美匹配你要的需求——让低优先级的numbers可观察对象,要么等高优先级的letters凑够2个一组再发射,要么超时后直接发射自身数据。
核心实现代码
import { zip, of, timer } from 'rxjs'; import { timeoutWith, mergeMap } from 'rxjs/operators'; // 模拟高优先级的letters流 const letters = of('a', 'b', 'c', 'd', 'e'); // 模拟低优先级的numbers流 const numbers = of(1, 2, 3, 4); const prioritizedQueue = numbers.pipe( mergeMap(num => // 等待letters凑出2个一组 zip(letters, letters).pipe( // 超时后直接发射当前的num timeoutWith(1000, of(num)) ) ) ); // 订阅测试 prioritizedQueue.subscribe(val => console.log(val));
逻辑拆解
- 用
mergeMap包裹每个numbers的发射值,让单个数值进入独立的等待逻辑; - 内部的
zip(letters, letters)会一直等待letters发射出2个值,实现“凑够2个一组”的高优先级需求; timeoutWith是关键:如果在设定的超时时间(示例中是1000毫秒)内,letters没凑够2个,就直接跳过等待,发射当前的低优先级数值,避免低优先级流被无限阻塞。
额外优化提示
- 你可以根据业务场景灵活调整
timeoutWith的超时时长; - 如果
letters是异步流(比如来自用户事件、API请求),这个逻辑依然适用,zip会自动等待两个值的到来; - 要是需要处理
letters最后剩余的单个值,可以额外结合bufferCount(2, 1)这类操作符,覆盖边界场景。
内容的提问来源于stack exchange,提问作者Cokorda Raka
相关产品推荐
相关产品推荐

