Limit concurrent blocks processed

This commit is contained in:
Samy Pessé
2015-02-13 00:04:27 +01:00
parent 57d5c23eaf
commit 7c83869869
2 changed files with 65 additions and 10 deletions
+14 -10
View File
@@ -6,6 +6,7 @@ var nunjucks = require("nunjucks");
var git = require("./utils/git");
var stringUtils = require("./utils/string");
var fs = require("./utils/fs");
var batch = require("./utils/batch");
var pkg = require("../package.json");
@@ -368,16 +369,19 @@ TemplateEngine.prototype.postProcess = function(content) {
return Q(content)
.then(that.replaceBlocks)
.then(function(content) {
return Q.all(_.map(that.blocks, function(blk, blkId) {
return Q()
.then(function() {
if (!blk.post) return Q();
return blk.post();
})
.then(function() {
delete that.blocks[blkId];
});
}))
return batch.execEach(that.blocks, {
max: 20,
fn: function(blk, blkId) {
return Q()
.then(function() {
if (!blk.post) return Q();
return blk.post();
})
.then(function() {
delete that.blocks[blkId];
});
}
})
.thenResolve(content);
});
};
+51
View File
@@ -0,0 +1,51 @@
var Q = require("q");
var _ = require("lodash");
// Execute a method for all element
function execEach(items, options) {
var concurrents = 0, d = Q.defer(), pending = [];
options = _.defaults(options || {}, {
max: 100,
fn: function(item) {}
});
function startItem(item, i) {
if (concurrents >= options.max) {
pending.push([item, i]);
return;
}
concurrents++;
Q()
.then(function() {
return options.fn(item, i);
})
.then(function() {
concurrents--;
// Next pending
var next = pending.shift();
if (concurrents == 0 && !next) {
d.resolve();
} else if (next) {
startItem.apply(null, next);
}
})
.fail(function(err) {
pending = [];
d.reject(err);
})
}
_.each(items, startItem);
return d.promise;
}
module.exports = {
execEach: execEach
};