【问题标题】:Lossless rate-limiting in RxJS with queue clearingRxJS 中使用队列清除的无损速率限制
【发布时间】:2016-10-27 13:47:10
【问题描述】:

在 rxjs5 中,我正在尝试实现 Throttler 类。

import Rx from 'rxjs/rx';

export default class Throttler {
  constructor(interval) {
    this.timeouts = [];
    this.incomingActions = new Rx.Subject();
    this.incomingActions
      .concatMap(action => Rx.Observable.just(action).delay(interval / 2))
      .subscribe(action => action());
  }

  clear() {
    // How do I do this?
  }

  do(action) {
    this.incomingActions.next(action);
  }
}

以下不变量必须成立:

  • 传递给do 的每个动作都被添加到动作队列中

  • 动作队列按构造函数参数确定的固定间隔按顺序处理

  • 可以使用clear() 清除操作队列。

如上所示,我当前的实现处理固定间隔,但我不知道如何清除队列。它还有一个问题,即使队列为空,所有操作都会延迟interval / 2ms。

附:我描述不变量的方式很容易映射到使用 setInterval 和数组作为队列的实现,但我想知道如何使用 Rx 来做到这一点。

【问题讨论】:

    标签: rxjs rxjs5


    【解决方案1】:

    这似乎不是默认Subject 类的好地方。由于您列出的原因,使用您自己的子类扩展它会更好。

    但是,在您的情况下,我会尝试使用一些索引来识别 .do(action) 方法的每个操作,并在 subscribe() 之前添加 .filter() 运算符,以便能够通过检查一些数组的索引来取消特定操作被标记为已取消。由于您使用的是concatMap(),因此您知道操作将始终按照添加的顺序被调用。然后你想要的clear() 方法只会在数组中标记所有要取消的操作。

    您还可以在concatMap() 之后添加.do() 运算符,并使用一些累加器跟踪当前安排了多少操作。添加操作将导致scheduledAction++,而在.subscribe() 之前通过.do() 将导致scheduledAction--。然后,您可以使用此变量来决定是否要使用 .delay(interval / 2) 链接新操作。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-05-04
      • 2017-06-18
      • 1970-01-01
      • 1970-01-01
      • 2020-12-05
      • 1970-01-01
      相关资源
      最近更新 更多