Refactored batch service to avoid event noise, doing in callback way

This commit is contained in:
Daniel García Aubert 2016-01-08 15:47:59 +01:00
parent f9f52d2bd1
commit 20f00d58d9
7 changed files with 158 additions and 163 deletions

View File

@ -28,6 +28,7 @@ function JobController(metadataBackend, tableCache, statsd_client) {
this.statsd_client = statsd_client; this.statsd_client = statsd_client;
this.jobBackend = new JobBackend(metadataBackend, jobQueue, jobPublisher, userIndexer); this.jobBackend = new JobBackend(metadataBackend, jobQueue, jobPublisher, userIndexer);
this.userDatabaseMetadataService = new UserDatabaseMetadataService(metadataBackend); this.userDatabaseMetadataService = new UserDatabaseMetadataService(metadataBackend);
this.jobCanceller = new JobCanceller(this.metadataBackend, this.userDatabaseMetadataService, this.jobBackend);
} }
JobController.prototype.route = function (app) { JobController.prototype.route = function (app) {
@ -102,20 +103,16 @@ JobController.prototype.cancelJob = function (req, res) {
req.profiler.done('setDBAuth'); req.profiler.done('setDBAuth');
} }
var jobCanceller = new JobCanceller(self.metadataBackend, self.userDatabaseMetadataService);
jobCanceller.cancel(job_id) self.jobCanceller.cancel(job_id, function (err, job) {
.on('cancelled', function (job) { if (err) {
// job is cancelled but surelly jobRunner has not deal whith it yet and it's not saved return next(err);
job.status = 'cancelled'; }
next(null, { next(null, {
job: job, job: job,
host: userDatabase.host host: userDatabase.host
}); });
})
.on('error', function (err) {
next(err);
}); });
}, },
function handleResponse(err, result) { function handleResponse(err, result) {

View File

@ -27,7 +27,7 @@ Batch.prototype.start = function () {
// do forever, it does not cause a stack overflow // do forever, it does not cause a stack overflow
forever(function (next) { forever(function (next) {
self._consume(host, queue, next); self._consumeJobs(host, queue, next);
}, function (err) { }, function (err) {
self.jobQueuePool.remove(host); self.jobQueuePool.remove(host);
@ -44,7 +44,7 @@ Batch.prototype.stop = function () {
this.jobSubscriber.unsubscribe(); this.jobSubscriber.unsubscribe();
}; };
Batch.prototype._consume = function consume(host, queue, callback) { Batch.prototype._consumeJobs = function (host, queue, callback) {
var self = this; var self = this;
queue.dequeue(host, function (err, job_id) { queue.dequeue(host, function (err, job_id) {
@ -58,20 +58,14 @@ Batch.prototype._consume = function consume(host, queue, callback) {
return callback(emptyQueueError); return callback(emptyQueueError);
} }
self.jobRunner.run(job_id) self.jobRunner.run(job_id, function (err, job) {
.on('done', function (job) { if (err) {
console.log('Job %s done in %s', job_id, host); return callback(err);
self.emit('job:done', job.job_id); }
callback();
}) console.log('Job %s %s in %s', job_id, job.status, host);
.on('failed', function (job) { self.emit('job:' + job.status, job_id);
console.log('Job %s failed in %s', job_id, host);
self.emit('job:failed', job.job_id);
callback();
})
.on('error', function (err) {
console.error('Error in job %s due to:', job_id, err.message || err);
self.emit('job:failed', job_id);
callback(); callback();
}); });
}); });

View File

@ -7,25 +7,19 @@ var UserDatabaseMetadataService = require('./user_database_metadata_service');
var JobPublisher = require('./job_publisher'); var JobPublisher = require('./job_publisher');
var JobQueue = require('./job_queue'); var JobQueue = require('./job_queue');
var UserIndexer = require('./user_indexer'); var UserIndexer = require('./user_indexer');
var JobBackend = require('./job_backend');
var Batch = require('./batch'); var Batch = require('./batch');
module.exports = function batchFactory (metadataBackend) { module.exports = function batchFactory (metadataBackend) {
var jobSubscriber = new JobSubscriber(); var jobSubscriber = new JobSubscriber();
var jobQueuePool = new JobQueuePool(metadataBackend); var jobQueuePool = new JobQueuePool(metadataBackend);
var userDatabaseMetadataService = new UserDatabaseMetadataService(metadataBackend);
var jobPublisher = new JobPublisher(); var jobPublisher = new JobPublisher();
var jobQueue = new JobQueue(metadataBackend); var jobQueue = new JobQueue(metadataBackend);
var userIndexer = new UserIndexer(metadataBackend); var userIndexer = new UserIndexer(metadataBackend);
var jobBackend = new JobBackend(metadataBackend, jobQueue, jobPublisher, userIndexer);
var jobRunner = new JobRunner( var userDatabaseMetadataService = new UserDatabaseMetadataService(metadataBackend);
metadataBackend, var jobRunner = new JobRunner(jobBackend, userDatabaseMetadataService);
userDatabaseMetadataService,
jobPublisher,
jobQueue,
userIndexer
);
return new Batch(jobSubscriber, jobQueuePool, jobRunner); return new Batch(jobSubscriber, jobQueuePool, jobRunner);
}; };

View File

@ -1,21 +1,17 @@
'use strict'; 'use strict';
var util = require('util');
var EventEmitter = require('events').EventEmitter;
var uuid = require('node-uuid'); var uuid = require('node-uuid');
var queue = require('queue-async'); var queue = require('queue-async');
var JOBS_TTL_AFTER_RESOLUTION = 48 * 3600; var JOBS_TTL_IN_SECONDS = 48 * 3600;
function JobBackend(metadataBackend, jobQueueProducer, jobPublisher, userIndexer) { function JobBackend(metadataBackend, jobQueueProducer, jobPublisher, userIndexer) {
EventEmitter.call(this); this.db = 5;
this.redisPrefix = 'batch:jobs:';
this.metadataBackend = metadataBackend; this.metadataBackend = metadataBackend;
this.jobQueueProducer = jobQueueProducer; this.jobQueueProducer = jobQueueProducer;
this.jobPublisher = jobPublisher; this.jobPublisher = jobPublisher;
this.userIndexer = userIndexer; this.userIndexer = userIndexer;
this.db = 5;
this.redisPrefix = 'batch:jobs:';
} }
util.inherits(JobBackend, EventEmitter);
JobBackend.prototype.create = function (username, sql, host, callback) { JobBackend.prototype.create = function (username, sql, host, callback) {
var self = this; var self = this;
@ -188,92 +184,107 @@ JobBackend.prototype.get = function (job_id, callback) {
}); });
}; };
JobBackend.prototype.setRunning = function (job) { JobBackend.prototype.setRunning = function (job, callback) {
var self = this; var now = new Date().toISOString();
var redisParams = [ var redisParams = [
this.redisPrefix + job.job_id, this.redisPrefix + job.job_id,
'status', 'running', 'status', 'running',
'updated_at', new Date().toISOString() 'updated_at', now
]; ];
this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams, function (err) { this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams, function (err) {
if (err) { if (err) {
return self.emit('error', err); return callback(err);
} }
self.emit('running', job); job.status = 'running';
job.updated_at = now;
callback(null, job);
}); });
}; };
JobBackend.prototype.setDone = function (job) { JobBackend.prototype.setDone = function (job, callback) {
var self = this; var self = this;
var now = new Date().toISOString();
var redisKey = this.redisPrefix + job.job_id; var redisKey = this.redisPrefix + job.job_id;
var redisParams = [ var redisParams = [
redisKey, redisKey,
'status', 'done', 'status', 'done',
'updated_at', new Date().toISOString() 'updated_at', now
]; ];
this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams , function (err) { this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams , function (err) {
if (err) { if (err) {
return self.emit('error', err); return callback(err);
} }
self.metadataBackend.redisCmd(self.db, 'EXPIRE', [ redisKey, JOBS_TTL_AFTER_RESOLUTION ], function (err) { self.metadataBackend.redisCmd(self.db, 'EXPIRE', [ redisKey, JOBS_TTL_IN_SECONDS ], function (err) {
if (err) { if (err) {
return self.emit('error', err); return callback(err);
} }
self.emit('done', job); job.status = 'done';
job.updated_at = now;
callback(null, job);
}); });
}); });
}; };
JobBackend.prototype.setFailed = function (job, err) { JobBackend.prototype.setFailed = function (job, err, callback) {
var self = this; var self = this;
var now = new Date().toISOString();
var redisKey = this.redisPrefix + job.job_id; var redisKey = this.redisPrefix + job.job_id;
var redisParams = [ var redisParams = [
redisKey, redisKey,
'status', 'failed', 'status', 'failed',
'failed_reason', err.message, 'failed_reason', err.message,
'updated_at', new Date().toISOString() 'updated_at', now
]; ];
this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams , function (err) { this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams , function (err) {
if (err) { if (err) {
return self.emit('error', err); return callback(err);
} }
self.metadataBackend.redisCmd(self.db, 'EXPIRE', [ redisKey, JOBS_TTL_AFTER_RESOLUTION ], function (err) { self.metadataBackend.redisCmd(self.db, 'EXPIRE', [ redisKey, JOBS_TTL_IN_SECONDS ], function (err) {
if (err) { if (err) {
return self.emit('error', err); return callback(err);
} }
self.emit('failed', job); job.status = 'failed';
job.updated_at = now;
callback(null, job);
}); });
}); });
}; };
JobBackend.prototype.setCancelled = function (job) { JobBackend.prototype.setCancelled = function (job, callback) {
var self = this; var self = this;
var now = new Date().toISOString();
var redisKey = this.redisPrefix + job.job_id; var redisKey = this.redisPrefix + job.job_id;
var redisParams = [ var redisParams = [
redisKey, redisKey,
'status', 'cancelled', 'status', 'cancelled',
'updated_at', new Date().toISOString() 'updated_at', now
]; ];
this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams , function (err) { this.metadataBackend.redisCmd(this.db, 'HMSET', redisParams , function (err) {
if (err) { if (err) {
return self.emit('error', err); return callback(err);
} }
self.metadataBackend.redisCmd(self.db, 'EXPIRE', [ redisKey, JOBS_TTL_AFTER_RESOLUTION ], function (err) { self.metadataBackend.redisCmd(self.db, 'EXPIRE', [ redisKey, JOBS_TTL_IN_SECONDS ], function (err) {
if (err) { if (err) {
return self.emit('error', err); return callback(err);
} }
self.emit('cancelled', job); job.status = 'cancelled';
job.updated_at = now;
callback(null, job);
}); });
}); });

View File

@ -1,54 +1,51 @@
'use strict'; 'use strict';
var JobBackend = require('./job_backend');
var PSQL = require('cartodb-psql'); var PSQL = require('cartodb-psql');
var JobPublisher = require('./job_publisher');
var JobQueue = require('./job_queue');
var UserIndexer = require('./user_indexer');
function JobCanceller(metadataBackend, userDatabaseMetadataService) { function JobCanceller(metadataBackend, userDatabaseMetadataService, jobBackend) {
this.metadataBackend = metadataBackend; this.metadataBackend = metadataBackend;
this.userDatabaseMetadataService = userDatabaseMetadataService; this.userDatabaseMetadataService = userDatabaseMetadataService;
this.jobBackend = jobBackend;
} }
JobCanceller.prototype.cancel = function (job_id) { JobCanceller.prototype.cancel = function (job_id, callback) {
var self = this; var self = this;
var jobQueue = new JobQueue(this.metadataBackend);
var jobPublisher = new JobPublisher();
var userIndexer = new UserIndexer(this.metadataBackend);
var jobBackend = new JobBackend(this.metadataBackend, jobQueue, jobPublisher, userIndexer);
jobBackend.get(job_id, function (err, job) { self.jobBackend.get(job_id, function (err, job) {
if (err) { if (err) {
return jobBackend.emit('error', err); return callback(err);
} }
if (job.status === 'pending') { if (job.status === 'pending') {
return jobBackend.setCancelled(job); return self.jobBackend.setCancelled(job, callback);
} }
if (job.status !== 'running') { if (job.status !== 'running') {
return jobBackend.emit('error', new Error('Job is ' + job.status + ' nothing to do')); return callback(new Error('Job is ' + job.status + ' nothing to do'));
} }
self.userDatabaseMetadataService.getUserMetadata(job.user, function (err, userDatabaseMetadata) { self.userDatabaseMetadataService.getUserMetadata(job.user, function (err, userDatabaseMetadata) {
if (err) { if (err) {
return jobBackend.emit('error', err); return callback(err);
} }
var pg = new PSQL(userDatabaseMetadata, {}, { destroyOnError: true }); self._query(job, userDatabaseMetadata, callback);
});
});
};
var getPIDQuery = 'SELECT pid FROM pg_stat_activity WHERE query = \'' + JobCanceller.prototype._query = function (job, userDatabaseMetadata, callback) {
job.query + var pg = new PSQL(userDatabaseMetadata, {}, { destroyOnError: true });
var getPIDQuery = 'SELECT pid FROM pg_stat_activity WHERE query = \'' + job.query +
' /* ' + job.job_id + ' */\''; ' /* ' + job.job_id + ' */\'';
pg.query(getPIDQuery, function(err, result) { pg.query(getPIDQuery, function(err, result) {
if(err) { if(err) {
return jobBackend.emit('error', err); return callback(err);
} }
if (!result.rows[0] || !result.rows[0].pid) { if (!result.rows[0] || !result.rows[0].pid) {
return jobBackend.emit('error', new Error('Query not running currently')); return callback(new Error('Query not running currently'));
} }
var pid = result.rows[0].pid; var pid = result.rows[0].pid;
@ -56,22 +53,23 @@ JobCanceller.prototype.cancel = function (job_id) {
pg.query(cancelQuery, function (err, result) { pg.query(cancelQuery, function (err, result) {
if (err) { if (err) {
return jobBackend.emit('error', err); return callback(err);
} }
var isCancelled = result.rows[0].pg_cancel_backend; var isCancelled = result.rows[0].pg_cancel_backend;
if (!isCancelled) { if (!isCancelled) {
return jobBackend.emit('error', new Error('Query has not been cancelled')); return callback(new Error('Query has not been cancelled'));
} }
jobBackend.emit('cancelled', job); // JobRunner handles job status through the PG's client error handler (see JobRunner.run:48)
}); // Due to user needs feedback, this modifies to the current status and updated dat
}); job.updated_at = new Date().toISOString();
}); job.status = 'cancelled';
});
return jobBackend; callback(null, job);
});
});
}; };

View File

@ -1,74 +1,74 @@
'use strict'; 'use strict';
var JobBackend = require('./job_backend');
var PSQL = require('cartodb-psql'); var PSQL = require('cartodb-psql');
var QUERY_CANCELED = '57014'; var QUERY_CANCELED = '57014';
function JobRunner(metadataBackend, userDatabaseMetadataService, jobPublisher, jobQueue, userIndexer) { function JobRunner(jobBackend, userDatabaseMetadataService) {
this.metadataBackend = metadataBackend; this.jobBackend = jobBackend;
this.userDatabaseMetadataService = userDatabaseMetadataService; this.userDatabaseMetadataService = userDatabaseMetadataService;
this.jobPublisher = jobPublisher;
this.jobQueue = jobQueue;
this.userIndexer = userIndexer;
} }
JobRunner.prototype.run = function (job_id) { JobRunner.prototype.run = function (job_id, callback) {
var self = this; var self = this;
var jobBackend = new JobBackend(this.metadataBackend, this.jobQueue, this.jobPublisher, this.userIndexer); self.jobBackend.get(job_id, function (err, job) {
jobBackend.get(job_id, function (err, job) {
if (err) { if (err) {
return jobBackend.emit('error', err); return callback(err);
} }
if (job.status !== 'pending') { if (job.status !== 'pending') {
return jobBackend.emit('error', return callback(new Error('Cannot run job ' + job.job_id + ' due to its status is ' + job.status));
new Error('Cannot run job ' + job.job_id + ' due to its status is ' + job.status));
} }
self.userDatabaseMetadataService.getUserMetadata(job.user, function (err, userDatabaseMetadata) { self.userDatabaseMetadataService.getUserMetadata(job.user, function (err, userDatabaseMetadata) {
if (err) { if (err) {
return jobBackend.emit('error', err); return callback(err);
} }
self.jobBackend.setRunning(job, function (err, job) {
if (err) {
return callback(err);
}
self._query(job, userDatabaseMetadata, callback);
});
});
});
};
JobRunner.prototype._query = function (job, userDatabaseMetadata, callback) {
var self = this;
var pg = new PSQL(userDatabaseMetadata, {}, { destroyOnError: true }); var pg = new PSQL(userDatabaseMetadata, {}, { destroyOnError: true });
jobBackend.setRunning(job);
pg.query('SET statement_timeout=0', function (err) { pg.query('SET statement_timeout=0', function (err) {
if(err) { if(err) {
return jobBackend.setFailed(job, err); return self.jobBackend.setFailed(job, err, callback);
} }
// mark query to allow to users cancel their queries whether users request for it // mark query to allow to users cancel their queries whether users request for it
var sql = job.query + ' /* ' + job.job_id + ' */'; var sql = job.query + ' /* ' + job.job_id + ' */';
pg.eventedQuery(sql, function (err, query /* , queryCanceller */) { pg.eventedQuery(sql, function (err, query) {
if (err) { if (err) {
return jobBackend.setFailed(job, err); return self.jobBackend.setFailed(job, err, callback);
} }
query.on('error', function (err) { query.on('error', function (err) {
if (err.code === QUERY_CANCELED) { if (err.code === QUERY_CANCELED) {
return jobBackend.setCancelled(job); return self.jobBackend.setCancelled(job, callback);
} }
jobBackend.setFailed(job, err); self.jobBackend.setFailed(job, err, callback);
}); });
query.on('end', function (result) { query.on('end', function (result) {
if (result) { if (result) {
jobBackend.setDone(job); self.jobBackend.setDone(job, callback);
} }
}); });
}); });
}); });
});
});
return jobBackend;
}; };
module.exports = JobRunner; module.exports = JobRunner;

View File

@ -21,6 +21,7 @@ describe('batch module', function() {
var jobPublisher = new JobPublisher(); var jobPublisher = new JobPublisher();
var userIndexer = new UserIndexer(metadataBackend); var userIndexer = new UserIndexer(metadataBackend);
var jobBackend = new JobBackend(metadataBackend, jobQueue, jobPublisher, userIndexer); var jobBackend = new JobBackend(metadataBackend, jobQueue, jobPublisher, userIndexer);
var batch = new Batch(metadataBackend); var batch = new Batch(metadataBackend);
before(function () { before(function () {