diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-08-08 16:06:25 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-08-08 16:06:25 +0200 |
| commit | d6975fc22c7dac0e003d6befd66a4d15eb808462 (patch) | |
| tree | 80aebd9d19dc4128b1c527d3e00b83c66b848495 | |
| parent | 0bac5bdb83bbc596f223b708bc226d8b2787359a (diff) | |
| download | shrinker-d6975fc22c7dac0e003d6befd66a4d15eb808462.tar.gz shrinker-d6975fc22c7dac0e003d6befd66a4d15eb808462.zip | |
refactor tryProcess to be in line with other services
| -rw-r--r-- | worker.js | 20 |
1 files changed, 10 insertions, 10 deletions
@@ -126,9 +126,6 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { // wrap process() in message handler async function tryProcess() { - const msgs = new Set(msg_buffer); - msg_buffer.clear(); - profiler.done("buffer filled"); profiler = undefined; @@ -136,12 +133,20 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { telemetries: telemetry_data.size }); + const msgs = new Set(msg_buffer); + msg_buffer.clear(); + const telemetry_objects = new Set(telemetry_data); + telemetry_data.clear(); + // clean up to allow processor to accept while we wait for db clearTimeout(idle_timer); clearTimeout(load_timer); + idle_timer = undefined; + load_timer = undefined; + try { - await process(); + await process(telemetry_objects); } catch (err) { if (err instanceof Seq.TimeoutError) { // deadlocks / timeout @@ -176,12 +181,7 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { } // finish a whole batch - async function process() { - const telemetry_objects = new Set(telemetry_data); - idle_timer = undefined; - load_timer = undefined; - telemetry_data.clear(); - + async function process(telemetry_objects) { // aggregate record objects to do a bulk insert let participant_phase_records = []; |
