summaryrefslogtreecommitdiff
path: root/worker.js
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-24 23:34:40 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-24 23:34:40 +0200
commit94f025db5247bb7b88c64c1f8d3c2e7d0134033d (patch)
tree34cce7541a7d240569075ace25140890922da5e4 /worker.js
parent53632f22e9a67f3ef3d657f8ed67907b2037fa3a (diff)
downloadcruncher-94f025db5247bb7b88c64c1f8d3c2e7d0134033d.tar.gz
cruncher-94f025db5247bb7b88c64c1f8d3c2e7d0134033d.zip
rework cruncher
Diffstat (limited to 'worker.js')
-rw-r--r--worker.js315
1 files changed, 71 insertions, 244 deletions
diff --git a/worker.js b/worker.js
index 7b09a9d..c367111 100644
--- a/worker.js
+++ b/worker.js
@@ -2,21 +2,30 @@
/* jshint esnext:true */
"use strict";
+/*
+ * cruncher calculates global_point and player_point for stats.
+ * It listens to the queue `crunch` and expects a JSON with
+ * a string that is a `participant.api_id`.
+ * cruncher will update the sums of points
+ * and send a notification to web.
+ */
+
const amqp = require("amqplib"),
Promise = require("bluebird"),
+ fs = Promise.promisifyAll(require("fs")),
winston = require("winston"),
loggly = require("winston-loggly-bulk"),
Seq = require("sequelize"),
- sleep = require("sleep-promise"),
- hash = require("object-hash");
+ sleep = require("sleep-promise");
const RABBITMQ_URI = process.env.RABBITMQ_URI,
DATABASE_URI = process.env.DATABASE_URI,
LOGGLY_TOKEN = process.env.LOGGLY_TOKEN,
- // number of inserts in one statement
- CHUNKSIZE = parseInt(process.env.CHUNKSIZE) || 300,
- MAXCONNS = parseInt(process.env.MAXCONNS) || 10, // how many concurrent actions
- CRUNCHERS = parseInt(process.env.CRUNCHERS) || 4; // how many players to crunch concurrently
+ // size of connection pool
+ MAXCONNS = parseInt(process.env.MAXCONNS) || 3,
+ // number of participants to calculate at once
+ BATCHSIZE = parseInt(process.env.BATCHSIZE) || 1000,
+ LOAD_TIMEOUT = parseInt(process.env.LOAD_TIMEOUT) || 5; // s
const logger = new (winston.Logger)({
transports: [
@@ -36,16 +45,11 @@ if (LOGGLY_TOKEN)
json: true
});
-// helpers
-// split an array into arrays of max chunksize
-function* chunks(arr) {
- for (let c=0, len=arr.length; c<len; c+=CHUNKSIZE)
- yield arr.slice(c, c+CHUNKSIZE);
-}
(async () => {
- let seq, model, rabbit, ch;
+ let seq, rabbit, ch;
+ // connect to rabbit & db
while (true) {
try {
seq = new Seq(DATABASE_URI, {
@@ -61,246 +65,69 @@ function* chunks(arr) {
await sleep(5000);
}
}
- model = require("../orm/model")(seq, Seq);
- function cartesian(arr) {
- return Array.prototype.reduce.call(arr, function(a, b) {
- let ret = [];
- a.forEach(function(a) {
- b.forEach(function(b) {
- ret.push(a.concat([b]));
- });
- });
- return ret;
- }, [[]]);
- }
-
- // create a 3D array
- // [[ ["hero", "Vox"], ["hero", "Taka"], …], [ ["game_mode", "ranked"], … ], …]
- // (will not use "Vox" but the index instead)
- async function dimensions_for(dimensions, on) {
- let cache = [];
- await Promise.map(dimensions, async (d, idx) =>
- cache[idx] = (await d.findAll()).map(
- (o) => [d, o]).filter((t) =>
- t[1].get("dimension_on") == null
- || t[1].get("dimension_on") == on)
- );
- return cache;
- }
- const player_dimensions = await dimensions_for(
- [model.Series, model.Filter, model.Hero, model.Role,
- model.GameMode], "player"),
- global_dimensions = await dimensions_for(
- [model.Series, model.Filter, model.Hero, model.Role,
- model.GameMode, model.Skilltier, model.Build,
- model.Region], "global");
+ // load update SQL scripts; scripts use sequelize replacements
+ // for the `participant_api_id` array
+ const player_script = fs.readFileSync("crunch_player.sql", "utf8"),
+ global_script = fs.readFileSync("crunch_global.sql", "utf8");
- // return every possible [query, insert] combination
- function calculate_point(points) {
- return points.map((point) => {
- // Series and Filter are special
- let where_aggr = {},
- where_links = {};
- // create skeleton: where hero_id=$hero
- // for aggregation, use series as range
- // for links, use the id
- point.forEach((tuple) => {
- // [table name, table element]
- if (tuple[1].get("name") != "all") {
- // exclude series & filter, added below
- if (tuple[0].tableName == "filter") {
- // merge custom filters
- Object.assign(where_aggr,
- tuple[1].get("filter"));
- // series and skill_tier are ranged filters
- } else if (tuple[0].tableName == "series") {
- // use start < date < end comparison
- where_aggr.created_at = { $between: [
- tuple[1].get("start"),
- tuple[1].get("end")
- ] }
- } else if (tuple[0].tableName == "skill_tier") {
- where_aggr["$participant.skill_tier$"] = { $between: [
- tuple[1].get("start"),
- tuple[1].get("end")
- ] }
- // build is a special ranged filter
- } else if (tuple[0].tableName == "build") {
- // TODO!
- } else if (tuple[0].tableName == "region") {
- // not joined via id, joined via name
- where_aggr["$participant.shard_id$"] = tuple[1].get("name");
- // most filters are directly as $filter_id on participant
- } else where_aggr["$participant." + tuple[0].tableName + "_id$"] =
- tuple[1].id
- }
- where_links[tuple[0].tableName + "_id"] = tuple[1].get("id");
- });
- return [where_aggr, where_links];
- });
- }
+ // fill a buffer and execute an SQL on a bigger (> 1o) batch
+ const participants_player = new Set(),
+ participants = new Set(),
+ // store the msgs that should be ACKed
+ buffer = new Set();
+ let timeout = undefined;
- // create an array with every possible combination
- // hero x game mode x …
- // Vox x casual x …
- // SAW x casual x …
- // …
- // Vox x ranked x …
- // SAW x ranked x …
- // …
- // Vox x ANY x …
- const player_points = calculate_point(cartesian(player_dimensions)),
- global_points = calculate_point(cartesian(global_dimensions));
+ // set maximum allowed number of unacked msgs
+ await ch.prefetch(BATCHSIZE);
+ ch.consume("crunch", (msg) => {
+ const api_id = msg.content.toString();
+ if (msg.type == "global")
+ participants.add(api_id);
+ // else exclusively add data to player, used for player refresh
+ participants_player.add(api_id);
+ buffer.add(msg);
+ if (timeout == undefined) timeout = setTimeout(crunch, LOAD_TIMEOUT*1000);
+ if (buffer.size >= BATCHSIZE) crunch();
+ }, { noAck: false });
- await ch.prefetch(CRUNCHERS);
+ // execute the scripts
+ async function crunch() {
+ const profiler = logger.startTimer();
+ logger.info("crunching");
- ch.consume("crunch", async (msg) => {
- const player_id = msg.content.toString();
- logger.info("working",
- { type: msg.properties.type, id: player_id });
+ // prevent async issues
+ const api_ids_player = [...participants_player],
+ api_ids = [...participants],
+ msgs = new Set(buffer);
+ participants.clear();
+ participants_player.clear();
+ buffer.clear();
+ clearTimeout(timeout);
+ timeout = undefined;
- let profiler = logger.startTimer();
- // service wide stats
- if (msg.properties.type == "global") {
- try {
- await calculate_global_point();
- logger.info("acking");
- await ch.ack(msg);
- // tell web
- await ch.publish("amq.topic", "global", new Buffer("points_update"));
- } catch (err) {
- logger.error("SQL error", err);
- await ch.nack(msg, false, true); // requeue
- }
- }
- // player stats
- if (msg.properties.type == "player") {
- try {
- await calculate_player_point(player_id);
- logger.info("acking");
- await ch.ack(msg);
- } catch (err) {
- logger.error("SQL error", err);
- await ch.nack(msg, false, true); // requeue
- }
- // tell web
- const player = await model.Player.findOne({
- where: { api_id: player_id },
- attributes: ["name"]
+ if (api_ids.length > 0)
+ await seq.query(global_script, {
+ replacements: { participant_api_ids: api_ids },
+ type: seq.QueryTypes.UPSERT
+ });
+ if (api_ids_player.length > 0)
+ await seq.query(player_script, {
+ replacements: { participant_api_ids: api_ids_player },
+ type: seq.QueryTypes.UPSERT
});
- if (player != null) {
- logger.info("updated player", { name: player.get("name") });
- await ch.publish("amq.topic", "player." + player.get("name"),
- new Buffer("points_update"));
- }
- }
- profiler.done("calculations for " +
- msg.properties.type + " " + player_id);
- }, { noAck: false });
-
- async function calculate_global_point() {
- logger.info("crunching global stats, this could take a while");
- await Promise.map(global_points, async (tuple, idx, len) => {
- const progress = Math.floor(100*100 * (1-idx/len)) / 100,
- where_aggr = tuple[0], where_links = tuple[1];
- // aggregate participant_stats with our condition
- let stats = await aggregate_stats(where_aggr);
- stats.updated_at = seq.fn("NOW");
- logger.info("inserting global stat", { progress: progress });
- Object.assign(stats, where_links);
- await model.GlobalPoint.upsert(stats);
- }, { concurrency: MAXCONNS });
- logger.info("committing");
- }
-
- async function calculate_player_point(player_api_id) {
- let point_records = [];
- logger.info("crunching player", { id: player_api_id });
- await Promise.map(player_points, async (tuple) => {
- const where_aggr = tuple[0],
- where_links = tuple[1];
- // make it player specific
- where_aggr["$participant.player_api_id$"] = player_api_id;
- where_aggr["final"] = true; // only end of match stats
- 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");
- Object.assign(stats, where_links);
- point_records.push(stats);
- }
- }, { concurrency: MAXCONNS });
+ // ack
+ await Promise.map(msgs, async (m) => await ch.ack(m));
+ // notify web
+ // TODO notify for player too
+ await ch.publish("amq.topic", "global", new Buffer("points_update"));
- logger.info("inserting into db");
- await seq.transaction({ autocommit: false }, async (transaction) => {
- await Promise.map(chunks(point_records), async (p_r) =>
- model.PlayerPoint.bulkCreate(p_r, {
- updateOnDuplicate: [], // all
- transaction: transaction
- }), { concurrency: MAXCONNS }
- )
- });
+ profiler.done("crunched");
}
+})();
- // return aggregated stats based on $where as WHERE clauses
- async function aggregate_stats(where) {
- // in literals: q -> column name, e -> function or string
- const q = (qry) => seq.dialect.QueryGenerator.quote(qry),
- e = (qry) => seq.dialect.QueryGenerator.escape(qry);
-
- // alternative for win rate
- //[ seq.literal(`${e(seq.fn("sum", seq.cast(seq.col("participant.winner"), "int") ))} / ${e(seq.fn("count", seq.col("participant.id")))}`), "win_rate" ]
-
- // short to sum a participant row as player stat with the same name
- const sum = (name) => [ seq.fn("sum", seq.col("participant_stats." + name)), name ];
+process.on("unhandledRejection", function(reason, promise) {
+ logger.error(reason);
+});
- const data = await model.ParticipantStats.findOne({
- where: where,
- attributes: [
- [ seq.fn("count", seq.col("participant.id")), "played" ],
- [ seq.fn("sum", seq.col("duration")), "time_spent" ],
- [ seq.fn("sum", seq.cast(seq.col("participant.winner"), "int") ), "wins" ],
- sum("kills"),
- sum("deaths"),
- sum("assists"),
- sum("minion_kills"),
- sum("jungle_kills"),
- sum("non_jungle_minion_kills"),
- sum("crystal_mine_captures"),
- sum("turret_captures"),
- sum("kda_ratio"),
- sum("kill_participation"),
- sum("impact_score"),
- sum("objective_score"),
- sum("damage_cp_score"),
- sum("damage_wp_score"),
- sum("sustain_score"),
- sum("farm_lane_score"),
- sum("kill_score"),
- sum("objective_lane_score"),
- sum("farm_jungle_score"),
- sum("peel_score"),
- sum("kill_assist_score"),
- sum("objective_jungle_score"),
- sum("vision_score"),
- sum("heal_score"),
- sum("assist_score"),
- sum("utility_score"),
- sum("synergy_score"),
- sum("build_score"),
- sum("offmeta_score"),
- sum("kraken_captures"),
- sum("gold")
- ],
- include: [ {
- model: model.Participant,
- as: "participant",
- attributes: []
- } ]
- });
- return data.dataValues;
- }
-})();