You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 07:25:18