【发布时间】:2020-02-02 17:22:12
【问题描述】:
假设我希望 observable 定期发出值,直到另一个 observable 发出。所以我可以使用timer 和takeUntil 来实现这一点。
但是然后我想处理每个发出的值并在某些条件变为真时停止(错误)发出。所以我写下一段代码:
const { timer, Subject } = 'rxjs';
const { takeUntil, map } = 'rxjs/operators';
const s = new Subject();
let count = 0;
function processItem(item) {
count++;
return count < 3;
}
const value = "MyValue";
timer(0, 1000).pipe(
takeUntil(s),
map(() => {
if (processItem(value)) {
console.log(`Processing done`);
return value;
} else {
s.next(true);
s.complete();
console.error(`Processing error`);
throw new Error(`Stop pipe`);
}
}),
)
但我没有收到错误,而是完成了我的 Observable。
仅当我注释掉 takeUntil(s) 运算符时,我才会收到错误消息。
看起来当管道运算符完成时,它的值不会立即发出,而是在管道的下一次“迭代”结束时被记住并发出,然后被新结果替换,依此类推。在我的情况下,下一次迭代,当应该发出错误时,takeUntil 会阻止它。问题是我对这个假设是否正确,如果是,为什么 rxjs 是这样设计的?
【问题讨论】:
标签: javascript rxjs