diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-20 11:47:42 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-20 11:47:42 +0200 |
| commit | 4fb48f248de4bc54ed76c50af8b80c528c04e027 (patch) | |
| tree | 57ee87f5fc557b3ae5191103ef0a11d22df50aea | |
| parent | ec256bbdab10f6c239c1849766d37b792ff10939 (diff) | |
| download | cruncher-4fb48f248de4bc54ed76c50af8b80c528c04e027.tar.gz cruncher-4fb48f248de4bc54ed76c50af8b80c528c04e027.zip | |
insert in chunks
| -rw-r--r-- | package.json | 1 | ||||
| -rw-r--r-- | worker.js | 32 |
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", @@ -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"); |
