summaryrefslogtreecommitdiff
path: root/worker.js
diff options
context:
space:
mode:
Diffstat (limited to 'worker.js')
-rw-r--r--worker.js32
1 files changed, 23 insertions, 9 deletions
diff --git a/worker.js b/worker.js
index 4955592..65d5d45 100644
--- a/worker.js
+++ b/worker.js
@@ -3,6 +3,7 @@
"use strict";
const amqp = require("amqplib"),
+ Promise = require("bluebird"),
winston = require("winston"),
Seq = require("sequelize"),
sleep = require("sleep-promise"),
@@ -10,6 +11,8 @@ const amqp = require("amqplib"),
const RABBITMQ_URI = process.env.RABBITMQ_URI,
DATABASE_URI = process.env.DATABASE_URI,
+ // number of inserts in one statement
+ CHUNKSIZE = parseInt(process.env.CHUNKSIZE) || 100,
CRUNCHERS = process.env.CRUNCHERS || 4; // how many players to crunch concurrently
const logger = new (winston.Logger)({
@@ -22,6 +25,13 @@ const logger = new (winston.Logger)({
]
});
+// 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;
@@ -154,14 +164,18 @@ const logger = new (winston.Logger)({
logger.info("inserting into db");
await seq.transaction({ autocommit: false }, (transaction) => {
return Promise.all([
- model.PlayerPoint.bulkCreate(player_records, {
- updateOnDuplicate: [], // all
- transaction: transaction
- }),
- model.GlobalPoint.bulkCreate(global_records, {
- updateOnDuplicate: [],
- transaction: transaction
- })
+ Promise.map(chunks(player_records), async (p_r) =>
+ model.PlayerPoint.bulkCreate(p_r, {
+ updateOnDuplicate: [], // all
+ transaction: transaction
+ })
+ ),
+ Promise.map(chunks(global_records), async (g_r) =>
+ model.GlobalPoint.bulkCreate(g_r, {
+ updateOnDuplicate: [],
+ transaction: transaction
+ })
+ )
]);
});
logger.info("acking");
@@ -169,7 +183,7 @@ const logger = new (winston.Logger)({
} catch (err) {
// TODO
logger.error("SQL error: %s, %j, %s",
- err.name, err.errors, err.parent.sql);
+ err.name, err.errors, err.parent);
await ch.nack(msg, false, true); // requeue
}
transaction_profiler.done("database transaction");