【问题标题】:Limit concurrency of promise being run限制正在运行的承诺的并发性
【发布时间】:2016-12-11 05:58:15
【问题描述】:

我正在寻找一个 Promise 函数包装器,它可以在给定的 Promise 运行时限制/节流,以便在给定的时间只运行一定数量的 Promise。

在下面的情况下,delayPromise 不应同时运行,它们都应按先到先得的顺序一次运行一个。

import Promise from 'bluebird'

function _delayPromise (seconds, str) {
  console.log(str)
  return Promise.delay(seconds)
}

let delayPromise = limitConcurrency(_delayPromise, 1)

async function a() {
  await delayPromise(100, "a:a")
  await delayPromise(100, "a:b")
  await delayPromise(100, "a:c")
}

async function b() {
  await delayPromise(100, "b:a")
  await delayPromise(100, "b:b")
  await delayPromise(100, "b:c")
}

a().then(() => console.log('done'))

b().then(() => console.log('done'))

关于如何设置这样的队列有什么想法吗?

我有一个来自美妙Benjamin Gruenbaum 的“去抖”功能。我需要修改它以根据它自己的执行而不是延迟来限制承诺。

export function promiseDebounce (fn, delay, count) {
  let working = 0
  let queue = []
  function work () {
    if ((queue.length === 0) || (working === count)) return
    working++
    Promise.delay(delay).tap(function () { working-- }).then(work)
    var next = queue.shift()
    next[2](fn.apply(next[0], next[1]))
  }
  return function debounced () {
    var args = arguments
    return new Promise(function (resolve) {
      queue.push([this, args, resolve])
      if (working < count) work()
    }.bind(this))
  }
}

【问题讨论】:

  • async.js 和 queue.js 都支持可配置的并发。
  • 用于数组或管理给定函数及其实例的状态?
  • when a given promise is running so that only a set number of that promise is running at a given time - 代码有 6 个 Promise,每个 Promise 将只运行一次 - 同时运行,但给定的 Promise 只运行一次 - 这个问题充其量是措辞不佳
  • 这里要清楚一点。承诺不会“运行”。 Promise 是已经运行的异步操作结果的代理。
  • es6-promise-pool? promise-limit? cwait?不相关或未知?

标签: javascript node.js promise bluebird


【解决方案1】:

我不认为有任何库可以做到这一点,但实际上自己实现非常简单:

function queue(fn) { // limitConcurrency(fn, 1)
    var q = Promise.resolve();
    return function(x) {
        var p = q.then(function() {
            return fn(x);
        });
        q = p.reflect();
        return p;
    };
}

对于多个并发请求,它会有点棘手,但也可以完成。

function limitConcurrency(fn, n) {
    if (n == 1) return queue(fn); // optimisation
    var q = null;
    var active = [];
    function next(x) {
        return function() {
            var p = fn(x)
            active.push(p.reflect().then(function() {
                active.splice(active.indexOf(p), 1);
            })
            return [Promise.race(active), p];
        }
    }
    function fst(t) {
        return t[0];
    }
    function snd(t) {
        return t[1];
    }
    return function(x) {
        var put = next(x)
        if (active.length < n) {
            var r = put()
            q = fst(t);
            return snd(t);
        } else {
            var r = q.then(put);
            q = r.then(fst);
            return r.then(snd)
        }
    };
}

顺便说一句,您可能想看看actors modelCSP。他们可以简化处理这些事情,也有一些 JS 库供他们使用。

示例

import Promise from 'bluebird'

function sequential(fn) {
  var q = Promise.resolve();
  return (...args) => {
    const p = q.then(() => fn(...args))
    q = p.reflect()
    return p
  }
}

async function _delayPromise (seconds, str) {
  console.log(`${str} started`)
  await Promise.delay(seconds)
  console.log(`${str} ended`)
  return str
}

let delayPromise = sequential(_delayPromise)

async function a() {
  await delayPromise(100, "a:a")
  await delayPromise(200, "a:b")
  await delayPromise(300, "a:c")
}

async function b() {
  await delayPromise(400, "b:a")
  await delayPromise(500, "b:b")
  await delayPromise(600, "b:c")
}

a().then(() => console.log('done'))
b().then(() => console.log('done'))

// --> with sequential()

// $ babel-node test/t.js
// a:a started
// a:a ended
// b:a started
// b:a ended
// a:b started
// a:b ended
// b:b started
// b:b ended
// a:c started
// a:c ended
// b:c started
// done
// b:c ended
// done

// --> without calling sequential()

// $ babel-node test/t.js
// a:a started
// b:a started
// a:a ended
// a:b started
// a:b ended
// a:c started
// b:a ended
// b:b started
// a:c ended
// done
// b:b ended
// b:c started
// b:c ended
// done

【讨论】:

  • @ThomasReggi:所以这个例子有效,你还有什么想回答的吗?
  • 嗯,有人在不解释他们的理由的情况下大肆投票。
【解决方案2】:

使用节流承诺模块:

https://www.npmjs.com/package/throttled-promise

var ThrottledPromise = require('throttled-promise'),
    promises = [
        new ThrottledPromise(function(resolve, reject) { ... }),
        new ThrottledPromise(function(resolve, reject) { ... }),
        new ThrottledPromise(function(resolve, reject) { ... })
    ];

// Run promises, but only 2 parallel
ThrottledPromise.all(promises, 2)
.then( ... )
.catch( ... );

【讨论】:

  • 我看不出如何用这个库解决 OP 的问题,它会在不同的范围内创建承诺。你能实现ab函数或limitConcurrency吗?
  • 这是一个流行的库,从字面上看,是 op 的问题:“我正在寻找一个 promise 函数包装器,它可以在给定的 promise 运行时限制/节流,这样只有一组承诺在给定的时间运行。”如果你想限制 Promise,你应该使用 throttled-promise。 SO 用于回答问题,而不是按需编写代码。
  • 他正在寻找一个承诺返回函数的包装器,而不是一个包装的 Promise 构造函数。此外,该库似乎要求您使用ThrottledPromise.all,这不符合问题中所述的问题。请尝试实现a/b 示例,您可能会看到它是如何不起作用的——或者如果它起作用,那么该代码实际上回答了这个问题。请注意,SO 不是图书馆推荐服务,但需要对问题的实际答案,这确实涉及代码编写。
【解决方案3】:

我也有同样的问题。我写了一个库来实现它。代码是here。我创建了一个队列来保存所有的承诺。当您将一些 Promise 推送到队列时,队列头部的前几个 Promise 将被弹出并运行。一旦一个承诺完成,队列中的下一个承诺也将被弹出并运行。一次又一次,直到队列中没有Task。您可以查看代码以获取详细信息。希望这个库对您有所帮助。

【讨论】:

    【解决方案4】:

    优势

    • 您可以定义并发承诺的数量(接近同时的请求)
    • 一致的流程:一旦一个承诺解决,另一个请求开始,无需猜测服务器能力
    • 对数据阻塞具有鲁棒性,如果服务器停止片刻,它将等待,并且下一个任务不会启动,因为 允许时钟
    • 不要依赖第 3 方模块,它是 Vanila node.js

    第一件事是让 https 成为一个承诺,所以我们可以使用等待来检索数据(从示例中删除) 第二创建一个承诺调度程序,在任何承诺得到解决时提交另一个请求。 第三次拨打电话

    通过限制并发承诺的数量来限制请求

    const https = require('https')
    
    function httpRequest(method, path, body = null) {
      const reqOpt = { 
        method: method,
        path: path,
        hostname: 'dbase.ez-mn.net', 
        headers: {
          "Content-Type": "application/json",
          "Cache-Control": "no-cache"
        }
      }
      if (method == 'GET') reqOpt.path = path + '&max=20000'
      if (body) reqOpt.headers['Content-Length'] = Buffer.byteLength(body);
      return new Promise((resolve, reject) => {
      const clientRequest = https.request(reqOpt, incomingMessage => {
          let response = {
              statusCode: incomingMessage.statusCode,
              headers: incomingMessage.headers,
              body: []
          };
          let chunks = ""
          incomingMessage.on('data', chunk => { chunks += chunk; });
          incomingMessage.on('end', () => {
              if (chunks) {
                  try {
                      response.body = JSON.parse(chunks);
                  } catch (error) {
                      reject(error)
                  }
              }
              console.log(response)
              resolve(response);
          });
      });
      clientRequest.on('error', error => { reject(error); });
      if (body) { clientRequest.write(body)  }  
      clientRequest.end();
    
      });
    }
    
        const asyncLimit = (fn, n) => {
          const pendingPromises = new Set();
    
      return async function(...args) {
        while (pendingPromises.size >= n) {
          await Promise.race(pendingPromises);
        }
    
        const p = fn.apply(this, args);
        const r = p.catch(() => {});
        pendingPromises.add(r);
        await r;
        pendingPromises.delete(r);
        return p;
      };
    };
    // httpRequest is the function that we want to rate the amount of requests
    // in this case, we set 8 requests running while not blocking other tasks (concurrency)
    
    
    let ratedhttpRequest = asyncLimit(httpRequest, 8);
    
    // this is our datase and caller    
    let process = async () => {
      patchData=[
          {path: '/rest/slots/80973975078587', body:{score:3}},
          {path: '/rest/slots/809739750DFA95', body:{score:5}},
          {path: '/rest/slots/AE0973750DFA96', body:{score:5}}]
    
      for (let i = 0; i < patchData.length; i++) {
        ratedhttpRequest('PATCH', patchData[i].path,  patchData[i].body)
      }
      console.log('completed')
    }
    
    process() 
    

    【讨论】:

      【解决方案5】:

      串联运行异步进程的经典方式是使用async.jsasync.series()。如果你更喜欢基于 Promise 的代码,那么有一个 Promise 版本的 async.jsasync-q

      使用async-q,您可以再次使用series

      async.series([
          function(){return delayPromise(100, "a:a")},
          function(){return delayPromise(100, "a:b")},
          function(){return delayPromise(100, "a:c")}
      ])
      .then(function(){
          console.log(done);
      });
      

      同时运行其中两个将同时运行ab,但在每个中它们将是连续的:

      // these two will run concurrently but each will run
      // their array of functions sequentially:
      async.series(a_array).then(()=>console.log('a done'));
      async.series(b_array).then(()=>console.log('b done'));
      

      如果您想在a 之后运行b,则将其放入.then()

      async.series(a_array)
      .then(()=>{
          console.log('a done');
          return async.series(b_array);
      })
      .then(()=>{
          console.log('b done');
      });
      

      如果您不想按顺序运行每个进程,而是希望限制每个进程同时运行一定数量的进程,那么您可以使用parallelLimit()

      // Run two promises at a time:
      async.parallelLimit(a_array,2)
      .then(()=>console.log('done'));
      

      阅读 async-q 文档:https://github.com/dbushong/async-q/blob/master/READJSME.md

      【讨论】:

      • OP 不想知道如何在没有 async 关键字的情况下实现a/b,而是想知道如何实现limitConcurrency,以便在示例中delayPromisea 调用和b 不平行。
      • @Bergi:这就是我展示的。使用async.series() 将依次运行delayPromise
      • 没有。即使您同时调用ab,它们也应该是连续的。
      猜你喜欢
      • 2015-11-24
      • 2015-01-15
      • 2016-03-15
      • 1970-01-01
      • 2021-08-08
      • 2015-07-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多