寻求类似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
相关产品推荐
相关产品推荐

