如何使用RxJs缓存资源密集型进程,取消订阅时立即终止原进程?
问题描述
我有一个资源密集型进程,希望将其最新结果缓存一段时间。
在以下示例中,我创建了Observable getInterval,它代表该资源密集型进程。其结果通过getSharedInterval共享,最后一次发射的结果会被缓存5秒。
但我想要实现的是:当外部Observable取消订阅时,getInterval能立即完成,同时最新结果仍会被缓存5秒。这是否可行?
const getInterval = new Observable((source) => { const interval = setInterval(() => { source.next('next'); }, 1000); source.add(() => { clearInterval(interval); }); }).pipe( tap(() => console.log('I am running')), finalize(() => console.log('completed')) ); const getSharedInterval = getInterval.pipe( share({ resetOnRefCountZero: () => timer(5000), connector: () => new ReplaySubject(1), resetOnComplete: true, }) ); getSharedInterval .pipe( take(1), tap(() => console.log('I dont no want anymore')) ) .subscribe(console.log);
解决方案
可行。当前代码里的share配置会让原getInterval在订阅数归零时继续运行5秒才终止,不符合你的需求。可以通过手动管理订阅计数和终止信号的方式实现,以下是具体方案:
import { Observable, ReplaySubject, timer, takeUntil, tap, finalize } from 'rxjs'; const getInterval = new Observable((source) => { const interval = setInterval(() => { source.next('next'); }, 1000); source.add(() => { clearInterval(interval); }); }).pipe( tap(() => console.log('I am running')), finalize(() => console.log('completed')) ); const getSharedInterval = (() => { let refCount = 0; let cacheSubject = new ReplaySubject(1); let terminateSignal: ReplaySubject<void> | null = null; let sourceSubscription: any = null; let cacheClearTimer: any = null; return new Observable((subscriber) => { refCount++; // 首次订阅时启动原进程 if (refCount === 1) { terminateSignal = new ReplaySubject(1); sourceSubscription = getInterval .pipe(takeUntil(terminateSignal)) .subscribe({ next: val => cacheSubject.next(val), complete: () => cacheSubject.complete(), error: err => cacheSubject.error(err) }); } // 订阅缓存的结果 const sub = cacheSubject.subscribe(subscriber); return () => { sub.unsubscribe(); refCount--; // 最后一个订阅取消时,立即终止原进程并启动缓存过期定时器 if (refCount === 0) { terminateSignal?.next(); sourceSubscription?.unsubscribe(); clearTimeout(cacheClearTimer); cacheClearTimer = setTimeout(() => { cacheSubject = new ReplaySubject(1); }, 5000); } }; }); })(); // 测试代码 getSharedInterval .pipe( take(1), tap(() => console.log('I dont want anymore')) ) .subscribe(console.log);
核心逻辑:
- 手动维护订阅计数
refCount,跟踪当前订阅数量 - 首次订阅时启动原
getInterval进程,并用takeUntil绑定终止信号 - 最后一个订阅取消时,立即通过终止信号停止原进程,同时启动5秒定时器,到期后重置缓存Subject
- 5秒内的新订阅仍能获取缓存的最新结果,超时后新订阅会重新启动原进程
内容的提问来源于stack exchange,提问作者Yoolan
相关产品推荐
相关产品推荐

