【问题标题】:Avoid recursion with RxJS5 observables避免使用 RxJS5 可观察对象进行递归
【发布时间】: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


    【解决方案1】:

    有两个运算符可以帮助您:retryretryWhen。两者都在源 observable 上重新订阅,从而重试失败的操作。

    查看这个示例,我们有一个在第一次count 订阅时失败的可观察对象:

    let getObs = (count) => {
      return Rx.Observable.create((subs) => {
        console.log('Subscription count = ', count);
    
        if(count) {
          count--;
          subs.error("ERROR");
        } else {
          subs.next("SUCCESS");
          subs.complete();
        }
      
        return () => {};
      });
    };
    
    getObs(2).subscribe(console.log, console.log);
    getObs(2).retry(2).subscribe(console.log, console.log);
    getObs(3).retry(2).subscribe(console.log, console.log);
    &lt;script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.0.1/Rx.min.js"&gt;&lt;/script&gt;

    如你所见:

    • 如果我们按原样调用它,它将失败
    • 使用retry,我们可以,嗯,重试几次并取得成功响应
    • 如果 observable 失败太多 tomes retry 将放弃并沿链传播错误。

    您真正需要的是retryWhen,因为retry 会立即尝试再次执行操作。

    let getObs = (count) => {
      return Rx.Observable.create((subs) => {
        if(count) {
          count--;
          subs.error("ERROR");
        } else {
          subs.next("SUCCESS");
          subs.complete();
        }
      
        return () => {};
      });
    };
    
    getObs(2).retryWhen(errors => errors.delay(100))
      .subscribe(console.log, console.log);
    getObs(4).retryWhen(errors => errors.delay(100))
      .subscribe(console.log, console.log);
    &lt;script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.0.1/Rx.min.js"&gt;&lt;/script&gt;

    使用retryWhen 添加延迟很容易,但在尝试次数较多后强制它失败:

    let getObs = (count) => {
      return Rx.Observable.create((subs) => {
        if(count) {
          count--;
          subs.error("ERROR");
        } else {
          subs.next("SUCCESS");
          subs.complete();
        }
      
        return () => {};
      });
    };
    
    getObs(2)
      .retryWhen(errors => {
        return errors.delay(100).scan((errorCount, err) => {
          if(!errorCount) {
            throw err;
          }
          return --errorCount;
        }, 2);
      })
      .subscribe(console.log, console.log);
    
    getObs(4)
      .retryWhen(errors => {
        return errors.delay(100).scan((errorCount, err) => {
          if(!errorCount) {
            throw err;
          }
          return --errorCount;
        }, 2);
      })
      .subscribe(console.log, console.log);
    &lt;script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.0.1/Rx.min.js"&gt;&lt;/script&gt;

    最后,两次重试都会抛出错误,所以在获取锁时需要这样做:

        .flatMap(() => {
            return acquireLock(this)
                .switchMap(obj => {
                  if(obj.error) {
                    return Rx.Observable.throw(obj.error);
                  } else {
                    Rx.Observable.of(obj);
                  }
                })
                .retryWhen(...)
        })
    

    【讨论】:

      猜你喜欢
      • 2023-03-30
      • 1970-01-01
      • 2018-04-13
      • 2020-12-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多