Ditch inflateSync in favor of plain inflate; test iteration order; add d3-queue dep

This commit is contained in:
Andrew Pendleton
2016-04-14 18:49:34 -04:00
parent 4474634452
commit 405d987716
3 changed files with 72 additions and 39 deletions
+25 -16
View File
@@ -9,6 +9,7 @@ var sm = new (require('sphericalmercator'));
var sqlite3 = require('sqlite3'); var sqlite3 = require('sqlite3');
var tiletype = require('tiletype'); var tiletype = require('tiletype');
var ZXYStream = require('./zxystream'); var ZXYStream = require('./zxystream');
var queue = require('d3-queue').queue;
function noop(err) { function noop(err) {
if (err) throw err; if (err) throw err;
@@ -625,34 +626,42 @@ MBTiles.prototype.geocoderDataIterator = function(type) {
var doneSentinel = {}; var doneSentinel = {};
var _this = this; var _this = this;
var zlibQueue = queue(1);
var inflate = function(data, callback) {
zlibQueue.defer(function(cb) {
zlib.inflate(data, function(err, buf) {
callback(err, buf);
cb();
});
});
}
var sending = false; var sending = false;
var sendIfAvailable = function() { var sendIfAvailable = function() {
if (sending) return; if (sending) return;
sending = true; sending = true;
while (nextQueue.length && dataQueue.length) { while (nextQueue.length && dataQueue.length) {
var nextCb = nextQueue.shift(), data, cbValue; var nextCb = nextQueue.shift(), data;
if (dataQueue[0] == doneSentinel) { if (dataQueue[0] == doneSentinel) {
cbValue = { (function(nextCb) {
err: null, setImmediate(function() {
row: {value: undefined, done: true} nextCb(null, {value: undefined, done: true});
} });
})(nextCb);
} else { } else {
data = dataQueue.shift(); data = dataQueue.shift();
maybeRefillBuffer(); maybeRefillBuffer();
cbValue = {
err: data.err,
row: {value: {shard: data.row.shard, data: zlib.inflateSync(data.row.data)}, done: false}
};
}
// bind the callback and data now so that they don't change before the setImmediate
// callback is executed
(function(nextCb, cbValue) { (function(nextCb, cbValue) {
setImmediate(function() { inflate(data.row.data, function(err, buf) {
nextCb(cbValue.err, cbValue.row); nextCb(
}); data.err,
})(nextCb, cbValue); {value: {shard: data.row.shard, data: buf}, done: false}
);
})
})(nextCb, data);
}
} }
sending = false; sending = false;
+1
View File
@@ -26,6 +26,7 @@
"Konstantin Käfer <kkaefer>" "Konstantin Käfer <kkaefer>"
], ],
"dependencies": { "dependencies": {
"d3-queue": "~2.0.3",
"tiletype": "0.1.x", "tiletype": "0.1.x",
"sqlite3": "3.x", "sqlite3": "3.x",
"sphericalmercator": "~1.0.1" "sphericalmercator": "~1.0.1"
+30 -7
View File
@@ -2,6 +2,8 @@ var fs = require('fs');
var util = require('util'); var util = require('util');
var MBTiles = require('..'); var MBTiles = require('..');
var tape = require('tape'); var tape = require('tape');
var queue = require('d3-queue').queue;
var crypto = require('crypto');
var expected = { var expected = {
bounds: '-141.005548666451,41.6690855919108,-52.615930948992,83.1161164353916', bounds: '-141.005548666451,41.6690855919108,-52.615930948992,83.1161164353916',
@@ -61,8 +63,6 @@ tape('putGeocoderData', function(assert) {
to.startWriting(function(err) { to.startWriting(function(err) {
assert.ifError(err); assert.ifError(err);
to.putGeocoderData('term', 0, new Buffer('asdf'), function(err) { to.putGeocoderData('term', 0, new Buffer('asdf'), function(err) {
assert.ifError(err);
to.putGeocoderData('term', 1, new Buffer('ZZZZZ'), function(err) {
assert.ifError(err); assert.ifError(err);
to.stopWriting(function(err) { to.stopWriting(function(err) {
assert.ifError(err); assert.ifError(err);
@@ -75,21 +75,41 @@ tape('putGeocoderData', function(assert) {
}); });
}); });
}); });
});
tape('geocoderDataIterator', function(assert) { tape('geocoderDataIterator', function(assert) {
to.startWriting(function(err) {
assert.ifError(err);
// get a bunch of shards of different sizes and put them in in an arbitrary order
var shardIds = {};
var q = queue()
while (true) {
var id = Math.floor(Math.random() * Math.pow(2, 16));
if (shardIds[id]) continue;
shardIds[id] = 1;
q.defer(function(id, cb) {
to.putGeocoderData("term", id, crypto.randomBytes(Math.floor(Math.random()* 1024 * 1024)), cb);
}, id);
if (Object.keys(shardIds).length >= 50) break;
}
q.awaitAll(function() {
to.stopWriting(function(err) {
var it = to.geocoderDataIterator("term"); var it = to.geocoderDataIterator("term");
var data = []; var data = [];
var n = function(err, item) { var n = function(err, item) {
assert.ifError(err); assert.ifError(err);
if (item.done) { if (item.done) {
assert.equal(data.length, 2, "iterator produces two shards"); assert.equal(data.length, 51, "iterator produces 51 shards");
assert.equal(data[0].shard, 0); assert.equal(data[0].shard, 0);
assert.equal(data[0].data.toString(), "asdf"); assert.equal(data[0].data.toString(), "asdf", "first shard data is preserved");
assert.equal(data[1].shard, 1); var sorted = true;
assert.equal(data[1].data.toString(), "ZZZZZ"); for (var i = 1; i < data.length; i++) {
if (data[i - 1].shard >= data[i].shard) sorted = false;
}
assert.equal(sorted, true, "shards come back in order");
assert.end(); assert.end();
} else { } else {
@@ -98,6 +118,9 @@ tape('geocoderDataIterator', function(assert) {
} }
} }
it.asyncNext(n); it.asyncNext(n);
});
})
});
}) })
tape('getIndexableDocs', function(assert) { tape('getIndexableDocs', function(assert) {