【问题标题】:Avoiding recursion with RxJS5 in simple use case在简单用例中避免使用 RxJS5 进行递归
【发布时间】:2016-12-30 22:18:02
【问题描述】:

我试图弄清楚如何避免递归(如果可能的话)。我正在使用 RxJS 在队列上创建方法。如果队列不为空,则 drain 方法递归调用自身。 drain 方法目前被设计为从队列中一次删除一个项目,直到它为空。以下似乎适用于此目的,但我想知道是否可以避免递归调用排水。递归的问题在于,在递归完成之前,我可能无法从递归方法中“返回”任何项目(这是我的猜测)。

所以我的问题是:

我希望 drain() 的订阅者分别从队列中接收每个项目,而不是在方法完成递归时立即从队列中接收所有耗尽的项目。我怎样才能做到这一点?我可以使用递归方法完成此操作,还是只能使用非递归方法完成此操作?如果是后者,该怎么做?

此方法将耗尽队列,并在队列耗尽时停止尝试 empty 我们需要锁定,删除一个项目,然后解锁,每次,所以 最简单的方法是使用递归并重新调用 drain 方法 队列不为空

Queue.prototype.drain = function (opts) {

    opts = opts || {};
    const delay = opts.delay || 500;

    return this.init()
        .flatMap(() => {
            return acquireLock(this)
                .flatMap(obj => {
                    return acquireLockRetry(obj)
                });
        })
        .flatMap(obj => {
            return removeOneLine(this)
                .flatMap(l => {
                    return releaseLock(this, obj.id)
                        .map(obj => l);
                });
        })
        .flatMap(() => {
            return Rx.Observable.timer(delay)
                .flatMap(() => {
                    return this.drain()   /// <<< recurse
                        .takeUntil(this.isEmpty());  /// <<<< until
                });
        })
        .catch(e => {
            const force = !String(e.stack || e).match(/acquire lock timed out/);
            return releaseLock(this, force);
        });

};

//检查队列是否为空

Queue.prototype.isEmpty = function () {

    return this.init()
        .flatMap(() => {
            return acquireLock(this)
                .flatMap(obj => {
                    return acquireLockRetry(obj)
                })
        })
        .flatMap(obj => {
            return findFirstLine(this)
                .flatMap(l => {
                    return releaseLock(this, obj.id)
                        .map(obj => l);
                });
        })
        .filter(l => {
            // filter out any lines => only fire event if there is no line
            return !l;
        })
        .catch(e => {
            const force = !String(e.stack || e).match(/acquire lock timed out/);
            return releaseLock(this, force);
        });

};

【问题讨论】:

  • 递归是大多数时候要走的路。为什么不使用传递给drain()callback 来使用中断?
  • 是的,我对递归没有真正的问题,只是在这种情况下,我可能无法在它完成之前发出事件。我想触发一个传入的回调会起作用。我想他们也可以只传入一个 RxJS 观察者,而不是一个普通的 cb (?)。
  • acquireLockreleaseLockisEmptyfindFirstLine 是异步的呢?您是从文件或网络资源中读取,还是分布式的?如果这两件事都不是,我认为在这种情况下 Rx 可能是矫枉过正。
  • 所有这些方法都是异步的,我认为队列中没有一个方法是同步的。队列存在于磁盘上,而不是内存中,并且锁定机制是联网的,因为多个进程可能会从磁盘上的队列中读取。你可以在没有 RxJS 的情况下做到这一点,就像任何事情一样,但是 RxJS 可能会有所帮助,只是想学习它,这是一个很好的借口。队列是一组永无止境的值,因此是 Rx 的好用例。
  • 相信我,我并没有那么自虐,以至于我会尝试使用 RxJS 进行一堆同步调用:)

标签: javascript node.js recursion rxjs5


【解决方案1】:

一个可能的解决方案是这样的:

    const obs = new Rx.Subject();

    q.drain(obs).subscribe(function (v) {
        console.log('end result => ', v);
    });

    obs.subscribe(function (v) {
        console.log('next item that was drained => ', v);
    });

而drain方法就变成了:

Queue.prototype.drain = function (obs, opts) {

    opts = opts || {};

    const delay = opts.delay || 500;

    return this.init()
        .flatMap(() => {
            return acquireLock(this)
                .flatMap(obj => {
                    console.log(' drain lock id => ', obj.id);
                    return acquireLockRetry(obj)
                });
        })
        .flatMap(obj => {
            return removeOneLine(this)
                .flatMap(l => {
                    return releaseLock(this, obj.id)
                        .map(obj => {
                            obs.next(l);
                            return l;
                        });
                });
        })
        .flatMap(() => {
            return Rx.Observable.timer(500)
                .flatMap(() => {
                    return this.drain(obs, opts)
                        .takeUntil(this.isEmpty());
                });
        })
        .catch(e => {
            console.error('\n', ' => isEmpty() error => \n', e.stack || e);
            const force = !String(e.stack || e).match(/acquire lock timed out/);
            return releaseLock(this, force);
        });

};

这样,每次删除一个项目时,它都会触发一些东西,然后希望在最后,当队列完全耗尽时触发一个事件。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-01-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-20
    相关资源
    最近更新 更多