【发布时间】:2016-12-29 21:45:46
【问题描述】:
好的,所以我想避免使用可观察对象进行递归,使用外部和内部事件的组合,而不是调用相同的方法/函数。
现在我有这个:
Queue.prototype.deq = function (opts) {
opts = opts || {};
const noOfLines = opts.noOfLines || opts.count || 1;
const isConnect = opts.isConnect !== false;
let $dequeue = this.init()
.flatMap(() => {
return acquireLock(this)
.flatMap(obj => {
if(obj.error){
// if there is an error acquiring the lock we
// retry after 100 ms, which means using recursion
// because we call "this.deq()" again
return Rx.Observable.timer(100)
.flatMap(() => this.deq(opts));
}
else{
return makeGenericObservable()
.map(() => obj);
}
})
})
.flatMap(obj => {
return removeOneLine(this)
.map(l => ({l: l, id: obj.id}))
})
.flatMap(obj => {
return releaseLock(this, obj.id)
.map(() => obj.l)
})
.catch(e => {
console.error(e.stack || e);
return releaseLock(this);
});
if (isConnect) {
$dequeue = $dequeue.publish();
$dequeue.connect();
}
return $dequeue;
};
除了上面使用递归的方法(希望是正确的),我想使用一种更有事件的方法。如果从acquireLock函数传回的错误,我想每100ms重试一次,一旦成功我想停止,我不知道该怎么做,我很难测试它...... .这是对的吗?
Queue.prototype.deq = function (opts) {
// ....
let $dequeue = this.init()
.flatMap(() => {
return acquireLock(this)
.flatMap(obj => {
if(obj.error){
return Rx.Observable.interval(100)
.takeUntil(
acquireLock(this)
.filter(obj => !obj.error)
)
}
else{
// this is just an "empty" observable
// which immediately fires onNext()
return makeGenericObservable()
.map(() => obj);
}
})
})
// ...
return $dequeue;
};
有没有办法让它更简洁?我也想只重试5次左右。我的主要问题是 - 我怎样才能在间隔旁边创建一个计数,以便我每 100 毫秒重试一次,直到获得锁或计数达到 5?
我需要这样的东西:
.takeUntil(this or that)
也许我可以像这样简单地链接 takeUntils:
return Rx.Observable.interval(100)
.takeUntil(
acquireLock(this)
.filter(obj => !obj.error)
)
.takeUntil(++count < 5);
我可以这样做:
return Rx.Observable.interval(100)
.takeUntil(
acquireLock(this)
.filter(obj => !obj.error)
)
.takeUntil( Rx.Observable.timer(500));
但可能正在寻找不那么笨拙的东西
但我不知道在哪里(存储/跟踪)count 变量...
还希望使其更简洁,并可能检查其正确性。
我不得不说,如果这个东西有效,它是非常强大的编码结构。
【问题讨论】:
标签: recursion rxjs observable rxjs5