【问题标题】:Observable operators to "do" something after an amount of time has passed经过一段时间后,可观察的运算符“做”某事
【发布时间】:2018-09-27 19:33:26
【问题描述】:

在 rxjs 可观察链中,我如何在经过一定时间后访问可观察的当前值?本质上,我正在寻找类似tap operator 的东西,但只有在经过一定时间而没有看到可观察值的情况下才会执行。所以实际上它就像是点击和超时的组合。

我在想像下面这样的事情

observable$.pipe(
  first(x => x > 5),
  tapAfterTime(2000, x => console.log(x)),
  map(x => x + 1)
).subscribe(...);

这是一个虚构的例子,“tapAfterTime”函数不是真实的。但基本思想是,如果订阅后 2000 毫秒过去了,并且 observable 没有看到大于 5 的值,那么无论 observable 的当前值是什么,都执行 tapAfterTime 回调函数。如果我们在 2000 毫秒之前看到大于 5 的值,那么 tapAfterTime 回调将永远不会运行,但 map 函数将始终按预期运行。

是否有实现此目的的运算符或运算符的任何组合?

【问题讨论】:

标签: javascript angular typescript rxjs rxjs-pipeable-operators


【解决方案1】:

也许这真的太复杂了,也许值得一看。

这个想法是有 2 个不同的 observable,创建转换源 observable$,然后最终合并。

第一个 Observable,我们称之为 obsFilterAndMapped,是您进行过滤和映射的地方。

第二个 Observable,我们称它为 obsTapDelay,是一个 Observable,它会在第一个 Observable(即obsFilterAndMapped)发射时以一定的延迟触发一个新计时器 - 一旦 delayTime 被传递然后你执行你的tapAfterTime 动作 - 如果第一个 Observable 在 delayTime 被传递之前发出一个新值,那么一个新的计时器被创建。

这是实现这个想法的代码

const stop = new Subject<any>();
const obsShared = observable$.pipe(
    finalize(() => {
        console.log('STOP');
        stop.next();
        stop.complete()
    }),
    share()
);
const delayTime = 300;
const tapAfterTime = (value) => {
    console.log('tap with delay', value)
}; 

let valueEmitted;

const obsFilterAndMapped = obsShared.pipe(
    tap(val => valueEmitted = val),
    filter(i => i > 7),
    map(val => val + ' mapped')
);

const startTimer = merge(of('START'), obsFilterAndMapped);

const obsTapDelay = startTimer.pipe(
    switchMap(val => timer(delayTime).pipe(
        tap(() => tapAfterTime(valueEmitted)),
        switchMap(() => empty()),
    )),
    takeUntil(stop),
)

merge(obsFilterAndMapped, obsTapDelay)
.subscribe(console.log, null, () => console.log('completed'))

使用这种方法,您可以在源 observable$ 不发出任何内容的时间超过 delayTime 的任何时间执行您的 tapAfterTime 操作。换句话说,这不仅适用于observable$ 的首次发射,而且适用于其整个生命周期。

您可以使用以下输入测试此类代码

const obs1 = interval(100).pipe(
    take(10),
);
const obs2 = timer(2000, 100).pipe(
    take(10),
    map(val => val + 200),
);
const observable$ = merge(obs1, obs2);

我们甚至可以考虑在闭包中隐藏valueEmitted 全局变量,但是这会增加代码的复杂性并且可能不值得。

【讨论】:

    【解决方案2】:

    看看

    可能是这样的?

    --

    initiateTimer() {
        if (this.timerSub) {
            this.timerSub.unsubscribe();
        }
    
        this.timerSub = Rx.Observable.timer(2000)
            .take(1)
            .subscribe(this.showPopup.bind(this));
    }
    

    【讨论】:

      【解决方案3】:

      可能是这样的:

      let cancel;
      observable$.pipe(
        tap((x)=>clear=setTimeout(()=>console.log(x), 2000)),
         filter(x => x > 5),
         tap(x => clearTimeout(clear)),
         map(x => x + 1)
      );
      

      【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-10-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多