【问题标题】:Remap RXJS observable to a timer start, without stream interruption将可观察的 RXJS 重新映射到计时器启动,而不会中断流
【发布时间】:2021-10-08 06:58:33
【问题描述】:

在不中断原始流的情况下,RXJS 中将 observable 重新映射为计时器起始值的正确方法是什么?

obs.pipe(take(1000), startTimer())
   .subscribe(start => {
       // show how long it took to finish streaming 1000 values:
       const duration = Date.now() - start;
       console.log(duration); 
   });

我希望 startTimer 使用 start 重新映射到一次性订阅,但不会中断原始流,即在这种情况下,subscribe 将仅在所有 1000 个值都完成流式传输后触发。

如何实现这样的startTimer?它应该会产生一个一次性的Date.now() 值,以帮助测量完整的流持续时间。

或者我是否已经缺少标准解决方案?

更新-1

预期结果如下图,但无需将start创建为外部变量,而是使其成为流的一部分:

const start = Date.now();

obs.pipe(take(1000))
    .subscribe({
        complete() {
            const duration = Date.now() - start;
            console.log(duration);
        }
    });

我想让它成为流的一部分的原因是因为原始的 observable 和订阅者彼此非常分离,就像坐在不相关的源文件中一样。

附:或者,如果可能的话,最终发出duration 的解决方案也很好。

更新-2

最后,我使用了一个通用的drain 操作数,旨在排出一个可观察的流,然后在最后产生一个可观察的:

/**
 * Drains the source observable till it completes, and then posts a new value-observable.
 */
function drain<T>(value: T | Observable<T> | (() => T | Observable<T>)) {
    const v = () => {
        const a = typeof value === 'function' ? value.call(null) : value;
        return a instanceof Observable ? a : of(a);
    }
    return s => defer(() => s.pipe(filter(_ => false), c => concat(c, v()))) as Observable<T>;
}

使用这个操作数,我可以像这样重写startTimer

const startTimer = () => drain(Date.now);

【问题讨论】:

    标签: rxjs


    【解决方案1】:

    一些代码与你描述的方式几乎完全一样:

    function logRunTime<T>(prefix: string): MonoTypeOperatorFunction<T> {
      return s => defer(() => {
        const start = Date.now();
        return s.pipe(
          tap({
            complete: () => console.log(`${prefix}: ${Date.now() - start}ms`)
          })
        );
      });
    }
    
    interval(1000).pipe(
      take(10),
      logRunTime("Ten Seconds of Interval")
    ).subscribe(console.log);
    

    输出:

    0
    1
    2
    3
    4
    5
    6
    7
    8
    9
    Ten Seconds of Interval: 10014ms
    

    更新 1

    不要让原始 observable 停止发射值 [...] 我们只是不想要源值

    在我看来,要么你继续发出这些值,要么你没有。

    这是一个降低源排放的版本。
    这就是你所追求的吗?

    function reduceRunTime<T>(prefix: string): OperatorFunction<T, string> {
      return s => defer(() => {
        const start = Date.now();
        return s.pipe(
          filter(_ => false),
          c => concat(c, of(null)),
          map(_ => `${prefix}: ${Date.now() - start}ms`)
        );
      }) as Observable<string>;
    }
    
    interval(1000).pipe(
      take(10),
      reduceRunTime("Ten Seconds of Interval")
    ).subscribe(console.log);
    

    输出:

    Ten Seconds of Interval: 10013ms
    

    更新 2

    如果您不想要字符串,这将在可观察对象完成后发出开始时间。

    function startTimer() {
        return s => s.pipe(
            filter(_ => false),
            c => concat(c, of(Date.now()))
        ) as Observable<number>;
    }
    

    更新 3

    两种不同的行为

    我认为更新 2 可能清理得太多了。考虑这个例子:

    const timed$ = interval(500).pipe(
      take(5),
      startTimer()
    );
    
    const logDiff = (start: number) => console.log(Date.now() - start);
    
    timed$.subscribe(logDiff);
    
    setTimeout(() => {
      timed$.subscribe(logDiff);
    }, 1000);
    setTimeout(() => {
      timed$.subscribe(logDiff);
    }, 5000);
    

    输出:

    2521
    3507
    7511
    

    值得注意的是,因为 Observables 是惰性的(在订阅之前什么都不做),但是在创建 observable 时会调用 Date.now。您的 startTime 很可能在 observable 开始之前很久就设置好了。使 2.5s 的 observable 看起来需要 7.5s。

    使用 defer 解决了这个问题,因为它在订阅之前不会创建 observable。

    更新startTimer

    function startTimer() {
      return s => defer(() => s.pipe(
          filter(_ => false),
          c => concat(c, of(Date.now()))
      )) as Observable<number>;
    }
    

    以上示例的新输出:

    2521
    2507
    2511
    

    现在您可以做一些有趣的事情,例如运行相同的 observable 10 次,然后平均运行时间以更好地了解需要多长时间。

    const average = arr => arr.reduce( ( p, c ) => p + c, 0 ) / arr.length;
    
    concat(...Array.from(Array(10)).map(_ => timed$)).pipe(
      map(start => Date.now() - start),
      tap(console.log),
      toArray()
    ).subscribe(runs => console.log("Average Runtime: ", average(runs)));
    

    输出:

    2515
    2506
    2506
    2506
    2507
    2505
    2506
    2506
    2507
    2507
    Average Runtime: 2507.1
    

    【讨论】:

    • 此代码发出所有原始流值,但我只想要原始流完成后的计时器开始值。当然,duration 值也不错。
    • 那么“不中断原始流”是什么意思?
    • without interrupting the original stream = 不要让原始的 observable 停止发射值,因为我们希望我们的 subscribe 在所有值都发射后,我们只是不想要源值在我们的例子中,我们只想要startduration
    • @vitaly-t - 我认为清理后的版本中有一个错误,它只会发出开始时间。我已经在上面的另一个更新中证明了我的意思。这很容易解决。
    • @vitaly-t 我没有发现任何问题 :)
    猜你喜欢
    • 2019-03-29
    • 1970-01-01
    • 2020-05-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-27
    • 2018-01-28
    相关资源
    最近更新 更多