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

寻求类似audit/throttle的RxJS操作符:串行执行foo并仅处理最新队列值

解决方案

你需要的是串行执行高成本任务、仅保留任务执行期间最新输入的RxJS逻辑,可通过自定义操作符实现,完全匹配你的需求:

import { Observable, Observer, Subscription } from 'rxjs';

function serialLatest<T, R>(foo: (value: T) => Observable<R>) {
  return (source: Observable<T>) => {
    return new Observable((observer: Observer<R>) => {
      let isProcessing = false;
      let pendingValue: T | null = null;
      const sourceSub = source.subscribe({
        next: (value) => {
          if (!isProcessing) {
            isProcessing = true;
            executeFoo(value);
          } else {
            // 仅保留最新的待处理值,丢弃中间值
            pendingValue = value;
          }
        },
        error: (err) => observer.error(err),
        complete: () => {
          // 源流结束后,若还有未处理的最新值,执行完再关闭输出
          if (!isProcessing && pendingValue) {
            executeFoo(pendingValue);
          } else if (!isProcessing) {
            observer.complete();
          }
        }
      });

      function executeFoo(value: T) {
        const fooSub = foo(value).subscribe({
          next: (result) => observer.next(result),
          error: (err) => {
            isProcessing = false;
            observer.error(err);
          },
          complete: () => {
            isProcessing = false;
            // 若有最新待处理值,立即启动下一次任务
            if (pendingValue) {
              const nextValue = pendingValue;
              pendingValue = null;
              executeFoo(nextValue);
            } else if (sourceSub.closed) {
              observer.complete();
            }
          }
        });
      }

      return () => {
        sourceSub.unsubscribe();
      };
    });
  };
}

使用示例

import { of, interval, delay, concat } from 'rxjs';

// 模拟高成本后端计算:延迟2秒返回结果
const foo = (value: number) => of(value).pipe(delay(2000));

// 模拟源流:连续发射3个值,暂停1秒后再发射2个值
const source$ = interval(1000).pipe(take(3), concat(interval(1000).pipe(take(2), delay(1000))));

// 用自定义操作符处理流
source$.pipe(serialLatest(foo)).subscribe(console.log);

核心逻辑说明

  • 串行执行控制:通过isProcessing标记确保同一时间仅运行一个foo实例,避免并行消耗资源
  • 最新值保留:foo运行期间的所有输入仅更新pendingValue,不会堆积队列,最多只存一个待处理值
  • 自动续跑机制:每次foo完成后,若存在待处理的最新值,立即启动下一次计算,无需额外触发
  • 源流收尾处理:源流结束后,若还有未处理的最新请求,会执行完计算再关闭输出流

这个实现解决了audit操作符的局限性——它不会直接发射源值,而是将最新的待处理值再次传入foo执行,完美匹配你的场景:即便短时间内收到100次请求,也只会在当前计算完成后用最新的请求触发一次新计算,避免不必要的资源浪费。

内容的提问来源于stack exchange,提问作者striderhobbit

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 23:44:57