From 4fb48f248de4bc54ed76c50af8b80c528c04e027 Mon Sep 17 00:00:00 2001 From: schneefux Date: Thu, 20 Apr 2017 11:47:42 +0200 Subject: insert in chunks --- worker.js | 32 +++++++++++++++++++++++--------- 1 file changed, 23 insertions(+), 9 deletions(-) (limited to 'worker.js') 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 { 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"); -- cgit v1.3.1