From ea1884f81967d135646a547cfec64f5cb32e3449 Mon Sep 17 00:00:00 2001 From: Joe Date: Mon, 14 Jan 2019 07:32:58 -0600 Subject: [PATCH] Only send next chunk after the LogAggregator feedback arrived. --- .../casRpsLogExtractor.config.json.example | 1 - casRpsLogExtractor/casRpsLogExtractor.js | 25 ++++++++----------- 2 files changed, 10 insertions(+), 16 deletions(-) diff --git a/casRpsLogExtractor/casRpsLogExtractor.config.json.example b/casRpsLogExtractor/casRpsLogExtractor.config.json.example index 8ec4582..acd4ee8 100644 --- a/casRpsLogExtractor/casRpsLogExtractor.config.json.example +++ b/casRpsLogExtractor/casRpsLogExtractor.config.json.example @@ -2,7 +2,6 @@ "sourceTag": "CCERpsLog", "serverURI": "http://localhost:2222", "pollIntervalIdle": 10000, - "pollIntervalBusy": 1000, "UP_sqlConfig": { "user": "casdbuser", "password": "4KcLao2016!", diff --git a/casRpsLogExtractor/casRpsLogExtractor.js b/casRpsLogExtractor/casRpsLogExtractor.js index 74ceb0c..dfec0f4 100644 --- a/casRpsLogExtractor/casRpsLogExtractor.js +++ b/casRpsLogExtractor/casRpsLogExtractor.js @@ -45,8 +45,6 @@ let maxPKey = ""; let maxPKeyUnsent = ""; let allRows = []; -let currPollInterval = config.pollIntervalBusy; -let lastPollInterval; let theInterval; @@ -72,7 +70,7 @@ sql.connect(config.sqlConfig, err => { sqlRequest.query("select top " + maxResultsPerMessage + " * from rpslog where pkey > '" + maxPKey + "' order by pkey"); } - theInterval = setInterval(runNextPoll, currPollInterval); + runNextPoll(); sqlRequest.on('row', row => { outData = { sourceTag: config.sourceTag, processingType: processingType, data: row }; @@ -89,6 +87,7 @@ sql.connect(config.sqlConfig, err => { }); sqlRequest.on('done', result => { + clearInterval(theInterval); if (allRows.length>0){ request({ method: 'POST', uri: config.serverURI, body: JSON.stringify(allRows), json: true }).then(_ => { // count up the maxPKey only when the request was successful... @@ -96,21 +95,17 @@ sql.connect(config.sqlConfig, err => { fs.writeFileSync("casRpsLogExtractor.maxPKey.dat", maxPKey); }).catch(err=>{ logger.log("error","Sending Request to Service: " + err) - }); + }).finally(_=>{ + runNextPoll(); + }); process.stdout.write(">" + allRows.length + ">"); - if (allRows.length === maxResultsPerMessage) { - currPollInterval = config.pollIntervalBusy; - } else { - currPollInterval = config.pollIntervalIdle; - } + } else { if (currPollInterval !== lastPollInterval) { - clearInterval(theInterval); - theInterval = setInterval(runNextPoll, currPollInterval); - } - - allRows=[]; + theInterval = setInterval(runNextPoll, config.pollIntervalIdle); + } + process.stdout.write("."); } + allRows=[]; - process.stdout.write("."); }); })