Implement metadata queries with plain Promises
Remove usage of PhasedExecution This achives better query execution granularity and removes questionable usage of shared results object. It introduces a couple of behavior changes: * estimatedFeatureCount desn't ignore errors now * sample always uses estimatedFeatureCount,even if the actual count is also computed.
This commit is contained in:
@@ -1,5 +1,4 @@
|
||||
var queryUtils = require('../../utils/query-utils');
|
||||
const PhasedExecution = require('../../utils/phased-execution');
|
||||
const AggregationMapConfig = require('../../models/aggregation/aggregation-mapconfig');
|
||||
var SubstitutionTokens = require('../../utils/substitution-tokens');
|
||||
|
||||
@@ -35,15 +34,14 @@ MapnikLayerStats.prototype.is = function (type) {
|
||||
return this._types[type] ? this._types[type] : false;
|
||||
};
|
||||
|
||||
function queryPromise(dbConnection, query, setResults) {
|
||||
function queryPromise(dbConnection, query, adaptResults) {
|
||||
return new Promise(function(resolve, reject) {
|
||||
dbConnection.query(query, function (err, res) {
|
||||
err = setResults(err, res);
|
||||
if (err) {
|
||||
reject(err);
|
||||
}
|
||||
else {
|
||||
resolve();
|
||||
resolve(adaptResults(res));
|
||||
}
|
||||
});
|
||||
|
||||
@@ -60,15 +58,7 @@ function columnAggregations(field) {
|
||||
return [];
|
||||
}
|
||||
|
||||
/* Helper to add a task to the queries PhasedExecution
|
||||
* type can be either 'pre' (for pre-aggregation metadata) or 'post'
|
||||
* zoom is used only for post-aggregation metadata
|
||||
* query is a function that generates a metadata query from a data query
|
||||
* assign is a function to assign the results of the metadata query
|
||||
* if a assignDefault function is present, it will be used in case of error
|
||||
* during the query execution and any errors will be ignored
|
||||
*/
|
||||
function addStat(queries, ctx, type, zoom, query, assign, assignDefault=null) {
|
||||
function _getSQL(ctx, type, zoom, query) {
|
||||
let sql;
|
||||
if (type === 'pre') {
|
||||
sql = ctx.preQuery;
|
||||
@@ -76,157 +66,164 @@ function addStat(queries, ctx, type, zoom, query, assign, assignDefault=null) {
|
||||
else {
|
||||
sql = queryForZoom(ctx.aggrQuery, zoom);
|
||||
}
|
||||
sql = query(sql);
|
||||
queries.task(
|
||||
queryPromise(
|
||||
ctx.dbConnection,
|
||||
sql,
|
||||
(err, res) => {
|
||||
if (!err) {
|
||||
assign(res);
|
||||
}
|
||||
else if (assignDefault !== null) {
|
||||
assignDefault();
|
||||
return null;
|
||||
}
|
||||
return err;
|
||||
}
|
||||
)
|
||||
return query(sql);
|
||||
}
|
||||
|
||||
function _estimatedFeatureCount(ctx) {
|
||||
// TODO: restore -1 on errors behavior?
|
||||
return queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'pre', 0, queryUtils.getQueryRowEstimation),
|
||||
res => ({ estimatedFeatureCount: res.rows[0].rows })
|
||||
);
|
||||
}
|
||||
|
||||
function _estimatedFeatureCount(queries, ctx) {
|
||||
if (queries.results.estimatedFeatureCount === undefined) {
|
||||
// This is always computed; a default value of -1 is used in case of error
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'pre', 0,
|
||||
queryUtils.getQueryRowEstimation,
|
||||
res => queries.results.estimatedFeatureCount = res.rows[0].rows,
|
||||
() => queries.results.estimatedFeatureCount = -1
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
function _featureCount(queries, ctx) {
|
||||
function _featureCount(ctx) {
|
||||
if (ctx.metaOptions.featureCount) {
|
||||
// TODO: if ctx.metaOptions.columnStats we can combine this with column stats query
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'pre', 0,
|
||||
queryUtils.getQueryActualRowCount,
|
||||
res => queries.results.featureCount = res.rows[0].rows
|
||||
return queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'pre', 0, queryUtils.getQueryActualRowCount),
|
||||
res => ({ featureCount: res.rows[0].rows })
|
||||
);
|
||||
}
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
function _aggrFeatureCount(queries, ctx) {
|
||||
function _aggrFeatureCount(ctx) {
|
||||
if (ctx.metaOptions.hasOwnProperty('aggrFeatureCount')) {
|
||||
// We expect as zoom level as the value of aggrFeatureCount
|
||||
// TODO: it'd be nice to admit an array of zoom levels to
|
||||
// return metadata for multiple levels.
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'post', ctx.metaOptions.aggrFeatureCount || 0,
|
||||
queryUtils.getQueryActualRowCount,
|
||||
res => queries.results.aggrfeatureCount = res.rows[0].rows
|
||||
return queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'post', ctx.metaOptions.aggrFeatureCount || 0, queryUtils.getQueryActualRowCount),
|
||||
res => ({ aggrFeatureCount: res.rows[0].rows })
|
||||
);
|
||||
}
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
function _geometryType(queries, ctx) {
|
||||
if (ctx.metaOptions.geometryType && queries.results.geometryType === undefined) {
|
||||
function _geometryType(ctx) {
|
||||
if (ctx.metaOptions.geometryType) {
|
||||
const geometryColumn = AggregationMapConfig.getAggregationGeometryColumn();
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'pre', 0,
|
||||
sql => queryUtils.getQueryGeometryType(sql, geometryColumn),
|
||||
res => queries.results.geometryType = res.rows[0].geom_type
|
||||
return queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'pre', 0, sql => queryUtils.getQueryGeometryType(sql, geometryColumn)),
|
||||
res => ({ geometryType: res.rows[0].geom_type })
|
||||
);
|
||||
}
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
function _columns(queries, ctx) {
|
||||
function _columns(ctx) {
|
||||
if (ctx.metaOptions.columns || ctx.metaOptions.columnStats) {
|
||||
// note: post-aggregation columns are in layer.options.columns when aggregation is present
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'pre', 0,
|
||||
sql => queryUtils.getQueryLimited(sql, 0),
|
||||
res => queries.results.columns = formatResultFields(ctx.dbConnection, res.fields)
|
||||
return queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'pre', 0, sql => queryUtils.getQueryLimited(sql, 0)),
|
||||
res => formatResultFields(ctx.dbConnection, res.fields)
|
||||
);
|
||||
}
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
// combine a list of results merging the properties of all the objects
|
||||
// undefined results are admitted and ignored
|
||||
function mergeResults(results) {
|
||||
if (results) {
|
||||
if (results.length === 0) {
|
||||
return {};
|
||||
}
|
||||
return results.reduce((a, b) => {
|
||||
if (a === undefined) {
|
||||
return b;
|
||||
}
|
||||
if (b === undefined) {
|
||||
return a;
|
||||
}
|
||||
return Object.assign({}, a, b);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
function firstPhaseQueries(queries, ctx) {
|
||||
_estimatedFeatureCount(queries, ctx);
|
||||
_featureCount(queries, ctx);
|
||||
_aggrFeatureCount(queries, ctx);
|
||||
_geometryType(queries, ctx);
|
||||
_columns(queries, ctx);
|
||||
// deeper (1 level) combination of a list of objects:
|
||||
// mergeColumns([{ col1: { a: 1 }, col2: { a: 2 } }, { col1: { b: 3 } }]) => { col1: { a: 1, b: 3 }, col2: { a: 2 } }
|
||||
function mergeColumns(results) {
|
||||
if (results) {
|
||||
if (results.length === 0) {
|
||||
return {};
|
||||
}
|
||||
return results.reduce((a, b) => {
|
||||
let c = Object.assign({}, b || {}, a || {});
|
||||
Object.keys(c).forEach(key => {
|
||||
if (b.hasOwnProperty(key)) {
|
||||
c[key] = Object.assign(c[key], b[key]);
|
||||
}
|
||||
});
|
||||
return c;
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
function _sample(queries, ctx) {
|
||||
|
||||
function _sample(ctx, numRows) {
|
||||
if (ctx.metaOptions.sample) {
|
||||
const numRows = queries.results.featureCount === undefined ?
|
||||
queries.results.estimatedFeatureCount :
|
||||
queries.results.featureCount;
|
||||
const sampleProb = Math.min(ctx.metaOptions.sample / numRows, 1);
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'pre', 0,
|
||||
sql => queryUtils.getQuerySample(sql, sampleProb),
|
||||
res => queries.results.sample = res.rows
|
||||
return queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'pre', 0, sql => queryUtils.getQuerySample(sql, sampleProb)),
|
||||
res => ({ sample: res.rows })
|
||||
);
|
||||
}
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
function _columnStats(queries, ctx) {
|
||||
function _columnStats(ctx, columns) {
|
||||
if (!columns) {
|
||||
return Promise.resolve();
|
||||
}
|
||||
if (ctx.metaOptions.columnStats) {
|
||||
let queries = [];
|
||||
let aggr = [];
|
||||
Object.keys(queries.results.columns).forEach(name => {
|
||||
queries.push(new Promise(resolve => resolve(columns))); // add columns as first result
|
||||
Object.keys(columns).forEach(name => {
|
||||
aggr = aggr.concat(
|
||||
columnAggregations(queries.results.columns[name])
|
||||
columnAggregations(columns[name])
|
||||
.map(fn => `${fn}(${name}) AS ${name}_${fn}`)
|
||||
);
|
||||
if (queries.results.columns[name].type === 'string') {
|
||||
if (columns[name].type === 'string') {
|
||||
const topN = ctx.metaOptions.columnStats.topCategories || 1024;
|
||||
// TODO: ctx.metaOptions.columnStats.maxCategories
|
||||
// => use PG stats to dismiss columns with more distinct values
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'pre', 0,
|
||||
sql => queryUtils.getQueryTopCategories(sql, name, topN),
|
||||
res => queries.results.columns[name].categories = res.rows
|
||||
queries.push(
|
||||
queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'pre', 0, sql => queryUtils.getQueryTopCategories(sql, name, topN)),
|
||||
res => ({ [name]: { categories: res.rows } })
|
||||
)
|
||||
);
|
||||
}
|
||||
});
|
||||
addStat(
|
||||
queries,
|
||||
ctx,
|
||||
'pre', 0,
|
||||
sql => `SELECT ${aggr.join(',')} FROM (${sql}) AS __cdb_query`,
|
||||
res => {
|
||||
Object.keys(queries.results.columns).forEach(name => {
|
||||
columnAggregations(queries.results.columns[name]).forEach(fn => {
|
||||
queries.results.columns[name][fn] = res.rows[0][`${name}_${fn}`];
|
||||
queries.push(
|
||||
queryPromise(
|
||||
ctx.dbConnection,
|
||||
_getSQL(ctx, 'pre', 0, sql => `SELECT ${aggr.join(',')} FROM (${sql}) AS __cdb_query`),
|
||||
res => {
|
||||
let stats = {};
|
||||
Object.keys(columns).forEach(name => {
|
||||
stats[name] = {};
|
||||
columnAggregations(columns[name]).forEach(fn => {
|
||||
stats[name][fn] = res.rows[0][`${name}_${fn}`];
|
||||
});
|
||||
});
|
||||
});
|
||||
}
|
||||
return stats;
|
||||
}
|
||||
)
|
||||
);
|
||||
return Promise.all(queries).then(results => ({ columns: mergeColumns(results) }));
|
||||
}
|
||||
}
|
||||
|
||||
function secondPhaseQueries(queries, ctx) {
|
||||
_sample(queries, ctx);
|
||||
_columnStats(queries, ctx);
|
||||
return Promise.resolve({ columns });
|
||||
}
|
||||
|
||||
// This is adapted from SQL API:
|
||||
@@ -278,27 +275,33 @@ function (layer, dbConnection, callback) {
|
||||
let aggrQuery = layer.options.sql;
|
||||
let preQuery = layer.options.sql_raw || aggrQuery;
|
||||
|
||||
let context = {
|
||||
let ctx = {
|
||||
dbConnection,
|
||||
preQuery,
|
||||
aggrQuery,
|
||||
metaOptions: layer.options.metadata || {}
|
||||
};
|
||||
|
||||
let queries = new PhasedExecution();
|
||||
|
||||
// TODO: could save some queries if queryUtils.getAggregationMetadata() has been used and kept somewhere
|
||||
// we would set queries.results.estimatedFeatureCount and queries.results.geometryType
|
||||
// (if metaOptions.geometryType) from it.
|
||||
|
||||
// Queries will be executed in two phases, with results from the first phase needed
|
||||
// to define the queries of the second phase
|
||||
queries.phase(() => firstPhaseQueries(queries, context));
|
||||
queries.phase(() => secondPhaseQueries(queries, context));
|
||||
queries.run()
|
||||
.then(results => callback(null, results))
|
||||
.catch(error => callback(error));
|
||||
// TODO: compute _sample with _featureCount when available
|
||||
|
||||
Promise.all([
|
||||
_estimatedFeatureCount(ctx).then(
|
||||
({ estimatedFeatureCount }) => _sample(ctx, estimatedFeatureCount)
|
||||
.then(s => mergeResults([s, { estimatedFeatureCount }]))
|
||||
),
|
||||
_featureCount(ctx),
|
||||
_aggrFeatureCount(ctx),
|
||||
_geometryType(ctx),
|
||||
_columns(ctx).then(columns => _columnStats(ctx, columns))
|
||||
]).then(results => {
|
||||
callback(null, mergeResults(results));
|
||||
}).catch(error => {
|
||||
callback(error);
|
||||
});
|
||||
};
|
||||
|
||||
module.exports = MapnikLayerStats;
|
||||
|
||||
@@ -1,94 +0,0 @@
|
||||
/**
|
||||
* PhasedExecution handles the execution of async tasks (via Promises)
|
||||
* which have dependencies between them in a simplified manner.
|
||||
* Instead of using the complete task dependency graph, tasks
|
||||
* are organized into execution phases. So that tasks from a latter
|
||||
* phase will be initialized after tasks from previous phases have
|
||||
* finished.
|
||||
*
|
||||
* All tasks place their results in a shared object to make them
|
||||
* available to tasks of latter phases.
|
||||
*
|
||||
* Each phase is defined by a function that defines its tasks.
|
||||
*
|
||||
* Example:
|
||||
*
|
||||
* let p = new PhasedExecution();
|
||||
* // Define first phase with tasks 1 & 2
|
||||
* p.phase(() => {
|
||||
* console.log('At phase I', p.results);*
|
||||
* p.results.phase1 = 1
|
||||
* p.task(new Promise((resolve) => {
|
||||
* setTimeout( () => {
|
||||
* console.log('At task 1:', p.results);
|
||||
* p.results.task1 = 100;
|
||||
* resolve();
|
||||
* }, 400);
|
||||
* }));
|
||||
* p.task(new Promise((resolve) => {
|
||||
* setTimeout( () => {
|
||||
* console.log('At task 2:', p.results);
|
||||
* p.results.task2 = 200;
|
||||
* resolve();
|
||||
* }, 100);
|
||||
* }));
|
||||
* });
|
||||
* // Define second phase with tasks 3 & 4
|
||||
* p.phase(() => {
|
||||
* console.log('At phase II', p.results);
|
||||
* p.results.phase2 = 2
|
||||
* p.task(new Promise((resolve) => {
|
||||
* setTimeout( () => {
|
||||
* console.log('At task 3:', p.results);
|
||||
* p.results.task3 = 300;
|
||||
* resolve();
|
||||
* }, 50);
|
||||
* }));
|
||||
* p.task(new Promise((resolve) => {
|
||||
* setTimeout( () => {
|
||||
* console.log('At task 4:', p.results);
|
||||
* p.results.task4 = 400;
|
||||
* resolve();
|
||||
* }, 100);
|
||||
* }));
|
||||
* });
|
||||
* // Define third phase with task 5
|
||||
* p.phase(() => {
|
||||
* console.log('At phase III', p.results);
|
||||
* p.results.phase3 = 3
|
||||
* p.task(new Promise((resolve) => {
|
||||
* setTimeout( () => {
|
||||
* console.log('At task 5:', p.results);
|
||||
* p.results.task5 = 500;
|
||||
* resolve();
|
||||
* }, 50);
|
||||
* }));
|
||||
* });
|
||||
* // Execute all tasks
|
||||
* p.run().then((results) => {
|
||||
* console.log("RESULTS:", results);
|
||||
* }).catch((err) => {
|
||||
* console.log("ERROR:", error);
|
||||
* });
|
||||
*/
|
||||
module.exports = class PhasedExecution {
|
||||
constructor() {
|
||||
this.results = {};
|
||||
this.phases = [];
|
||||
}
|
||||
phase(phasegenerator) {
|
||||
this.phases.push(phasegenerator);
|
||||
}
|
||||
task(promise) {
|
||||
this.tasks.push(promise);
|
||||
}
|
||||
run() {
|
||||
this.tasks = [];
|
||||
let phase = this.phases.shift();
|
||||
if (phase) {
|
||||
phase(this);
|
||||
return Promise.all(this.tasks).then(() => this.run());
|
||||
}
|
||||
return this.results;
|
||||
}
|
||||
};
|
||||
@@ -2261,7 +2261,7 @@ describe('aggregation', function () {
|
||||
assert.equal(typeof body.metadata, 'object');
|
||||
assert.ok(Array.isArray(body.metadata.layers));
|
||||
assert.ok(body.metadata.layers[0].meta.aggregation.mvt);
|
||||
assert.equal(body.metadata.layers[0].meta.stats.aggrfeatureCount, 13);
|
||||
assert.equal(body.metadata.layers[0].meta.stats.aggrFeatureCount, 13);
|
||||
|
||||
done();
|
||||
});
|
||||
@@ -2302,7 +2302,7 @@ describe('aggregation', function () {
|
||||
assert.equal(typeof body.metadata, 'object');
|
||||
assert.ok(Array.isArray(body.metadata.layers));
|
||||
assert.ok(body.metadata.layers[0].meta.aggregation.mvt);
|
||||
assert.equal(body.metadata.layers[0].meta.stats.aggrfeatureCount, 9);
|
||||
assert.equal(body.metadata.layers[0].meta.stats.aggrFeatureCount, 9);
|
||||
|
||||
done();
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user