【发布时间】: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