summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-20 11:47:42 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-20 11:47:42 +0200
commit4fb48f248de4bc54ed76c50af8b80c528c04e027 (patch)
tree57ee87f5fc557b3ae5191103ef0a11d22df50aea
parentec256bbdab10f6c239c1849766d37b792ff10939 (diff)
downloadcruncher-4fb48f248de4bc54ed76c50af8b80c528c04e027.tar.gz
cruncher-4fb48f248de4bc54ed76c50af8b80c528c04e027.zip
insert in chunks
-rw-r--r--package.json1
-rw-r--r--worker.js32
2 files changed, 24 insertions, 9 deletions
diff --git a/package.json b/package.json
index 480fcb7..a3980f0 100644
--- a/package.json
+++ b/package.json
@@ -5,6 +5,7 @@
"main": "worker.js",
"dependencies": {
"amqplib": "^0.5.1",
+ "bluebird": "^3.5.0",
"mysql": "^2.13.0",
"object-hash": "^1.1.7",
"sequelize": "^3.30.4",
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");