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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 08:35:27