【问题标题】:Throttling asynchronous requests with RxJS使用 RxJS 限制异步请求
【发布时间】:2018-04-25 18:12:59
【问题描述】:

我需要向服务发送大量数据。每分钟允许我发送的项目有限制,每个请求我需要发送不超过 350 个项目,因此我将数据拆分为页面并尝试使用Observable 来限制请求:

const maxItemsPerRequest = 350;
const interval = Math.round(60 /*seconds*/ / maxPerMinute * 1000);
const pages = Math.ceil(totalItems / maxItemsPerRequest);

let observable = Observable.interval(interval).take(pages);

observable.subscribe(async page => {
  const items = await getItems(page, maxItemsPerRequest);
  await this.sendData(items);
});

getItems 和 sendData 可能需要一段时间才能完成,因此第二分钟可能会超出请求限制,例如(如果在第一分钟创建的请求需要超过 60 秒才能完成)。

如何确保订阅者在前一个请求完成后至少等待interval 毫秒,然后再发送新请求?

【问题讨论】:

    标签: javascript node.js rxjs


    【解决方案1】:

    基本上您不想在subscribe 中执行该操作,因为您无法控制它何时完成。相反,您想让它成为链的一部分,延迟将从上一个调用结束时开始。

    例如,您可以这样做:

    let observable = Observable.range(pages)
      .concatMap(page => getItems(page, maxItemsPerRequest)
        .concatMap(items => this.sendData(items))
        .delay(interval) // delay the emission from this Observable
      )
      .subscribe(console.log);
    

    【讨论】:

    • 相当不错。唯一的问题是你有任何延迟 sendData 给 plus 间隔。
    • 是的,这就是 OP 想要的:订阅者在发送新请求之前在上一个请求完成后至少等待间隔毫秒?
    • 确实,我在发布答案后注意到了这一点。尽管如此,他想要什么,他需要什么......
    • 但是,60 秒的等待时间也很长 - 延迟是累加的可能无关紧要。
    【解决方案2】:

    这是一个序列,它使用 Subject 在 sendData() 完成时发出信号,并根据该信号限制间隔。

    下面的演示有 500ms 的间隔,但是 sendData 导致了 1000ms 的延迟。如果将 sendData 的内部延迟更改为 100,则流以 500 毫秒的间隔继续。

    console.clear()
    
    const nextPage = new Rx.BehaviorSubject(true);
    
    const interval = 500 
    const pages = 8
    
    const getItems = () => Rx.Observable.of([1,2,3])
    
    const sendData = () => {
      return Rx.Observable.of('x')
        .delay(1000)
    }
    
    let observable = Rx.Observable.interval(interval)
      .throttle(i => nextPage)
      .take(pages)
      .concatMap(p => getItems())
      .concatMap(items => sendData(items))
      .do(x => nextPage.next())
    
    const start = Date.now()
    const show = (val) => console.log(val, Date.now() - start)
    observable.subscribe(show);
    <script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.5.4/Rx.js"></script>

    我的回答完全(嗯,不完全)不正确

    原来concatMap()是自我节流的,所以Subject是完全没有必要的。

    来自learn rxjs - concatMap

    concatMap 在前一个完成之前不会订阅下一个 observable

    这是没有主题的示例代码。

    console.clear()
    
    const interval = 500 
    const pages = 8
    
    const getItems = () => Rx.Observable.of([1,2,3])
    
    const sendData = () => {
      return Rx.Observable.of('x')
        .delay(1000)
    }
    
    let observable = Rx.Observable.interval(interval)
      .take(pages)
      .concatMap(p => getItems())
      .concatMap(items => sendData(items))
    
    const start = Date.now()
    const show = (val) => console.log(val, Date.now() - start)
    observable.subscribe(show);
    <script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.5.4/Rx.js"></script>

    【讨论】:

      猜你喜欢
      • 2016-06-11
      • 1970-01-01
      • 1970-01-01
      • 2017-11-15
      • 1970-01-01
      • 1970-01-01
      • 2014-08-22
      • 1970-01-01
      • 2017-06-13
      相关资源
      最近更新 更多