summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-13 13:40:54 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-13 13:40:54 +0200
commit0bb88bbf310749912ebcc13d92598d6b76f31539 (patch)
treec1ef7d1c366169b97f6514045ccfadeb9b2dba2f
parent6e92eac19aa19607d4734eeddb60fccdfd025535 (diff)
downloadcruncher-0bb88bbf310749912ebcc13d92598d6b76f31539.tar.gz
cruncher-0bb88bbf310749912ebcc13d92598d6b76f31539.zip
do not use batches
-rw-r--r--worker.js102
1 files changed, 37 insertions, 65 deletions
diff --git a/worker.js b/worker.js
index 2fa5fd2..74b971e 100644
--- a/worker.js
+++ b/worker.js
@@ -9,8 +9,7 @@ var amqp = require("amqplib"),
var RABBITMQ_URI = process.env.RABBITMQ_URI,
DATABASE_URI = process.env.DATABASE_URI,
- BATCHSIZE = parseInt(process.env.BATCHSIZE) || 5, // players + globals
- IDLE_TIMEOUT = parseFloat(process.env.IDLE_TIMEOUT) || 5000; // ms
+ PARALLELS = process.env.PARALLELS || 5; // how many players to crunch concurrently
(async () => {
let seq, model, rabbit, ch;
@@ -113,54 +112,29 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI,
let player_points = calculate_point(cartesian(player_dimensions)),
global_points = calculate_point(cartesian(global_dimensions));
- let queue = [],
- timer = undefined;
-
- await ch.prefetch(BATCHSIZE);
-
- ch.consume("crunch", (msg) => {
- queue.push(msg);
-
- // fill queue until batchsize or idle
- if (timer != undefined) clearTimeout(timer);
- timer = setTimeout(process, IDLE_TIMEOUT);
- if (queue.length == BATCHSIZE)
- process();
- }, { noAck: false });
-
- async function process() {
- console.log("processing batch", queue.length);
-
- // clean up to allow processor to accept while we wait for db
- clearTimeout(timer);
- timer = undefined;
- let msgs = queue;
- queue = [];
+ await ch.prefetch(PARALLELS);
+ ch.consume("crunch", async (msg) => {
let player_records = [],
global_records = [],
- players_done = [],
- global_done = false;
+ player_id = msg.content.toString();
- await Promise.all(msgs.map(async (msg) => {
- if (msg.properties.type == "global" && global_done == false) {
- let records = await calculate_global_point();
- if (records != undefined)
- global_records = global_records.concat(records);
- }
- if (msg.properties.type == "player") {
- let player_id = msg.content.toString();
- if (players_done.indexOf(player_id) == -1) {
- players_done.push(player_id);
- let records = await calculate_player_point(player_id);
- if (records != undefined)
- player_records = player_records.concat(records);
- }
- }
- }));
+ console.log("got work to do, let's go",
+ msg.properties.type, player_id);
+
+ if (msg.properties.type == "global" && global_done == false) {
+ let records = await calculate_global_point();
+ if (records != undefined)
+ global_records = global_records.concat(records);
+ }
+ if (msg.properties.type == "player") {
+ let records = await calculate_player_point(player_id);
+ if (records != undefined)
+ player_records = player_records.concat(records);
+ }
try {
- console.log("inserting batch into db");
+ console.log("inserting lots of stuff into db");
await seq.transaction({ autocommit: false }, (transaction) => {
return Promise.all([
model.PlayerPoint.bulkCreate(player_records, {
@@ -168,24 +142,33 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI,
transaction: transaction
}),
model.GlobalPoint.bulkCreate(global_records, {
- updateOnDuplicate: [], // all
+ updateOnDuplicate: [],
transaction: transaction
})
]);
});
- console.log("acking batch");
- await Promise.all(msgs.map((m) => ch.ack(m)) );
+ console.log("acking");
+ await ch.ack(msg);
} catch (err) {
// TODO
console.error(err);
- await Promise.all(msgs.map((m) => ch.nack(m, true)) ); // requeue
+ await ch.nack(msg, true); // requeue
}
- await Promise.all(players_done.map(async (p) =>
- await ch.publish("amq.topic", "player." + p.name, new Buffer("points_update")) ));
+ if (player_records.length > 0) {
+ let player = await model.Player.findOne({
+ where: { api_id: player_id },
+ attributes: ["name"]
+ });
+ if (player != null) {
+ console.log("finished updating player", player.get("name"));
+ await ch.publish("amq.topic", "player." + player.get("name"),
+ new Buffer("points_update"));
+ }
+ }
if (global_records.length > 0)
await ch.publish("amq.topic", "global", new Buffer("points_update"));
- }
+ }, { noAck: false });
async function calculate_global_point() {
let global_records = [];
@@ -207,30 +190,19 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI,
}
async function calculate_player_point(player_api_id) {
- let player = await model.Player.findOne(
- { where: { api_id: player_api_id } }),
- point_records = [];
- if (player == null) {
- console.error("player does not exist in db!?");
- return;
- }
- if (player.get("last_update") == null)
- // not enough data
- return;
- console.log("crunching player", player.name);
+ let point_records = [];
+ console.log("crunching player", player_api_id);
await Promise.all(player_points.map(async (tuple) => {
let where_aggr = tuple[0],
where_links = tuple[1];
// make it player specific
where_aggr["$participant.player_api_id$"] = player_api_id;
- where_links["player_api_id"] = player.api_id;
+ where_links["player_api_id"] = player_api_id;
// aggregate participant_stats with our condition
let stats = await aggregate_stats(where_aggr);
if (stats != undefined) {
stats.updated_at = seq.fn("NOW");
- stats.name = player.name;
- console.log("inserting stats for player", player.name);
Object.assign(stats, where_links);
point_records.push(stats);
}