diff --git a/casRpsLogExtractor/casRpsLogExtractor.config.json.example b/casRpsLogExtractor/casRpsLogExtractor.config.json.example index 129b5cf..8ec4582 100644 --- a/casRpsLogExtractor/casRpsLogExtractor.config.json.example +++ b/casRpsLogExtractor/casRpsLogExtractor.config.json.example @@ -1,7 +1,8 @@ { "sourceTag": "CCERpsLog", "serverURI": "http://localhost:2222", - "pollInterval": 10000, + "pollIntervalIdle": 10000, + "pollIntervalBusy": 1000, "UP_sqlConfig": { "user": "casdbuser", "password": "4KcLao2016!", diff --git a/casRpsLogExtractor/casRpsLogExtractor.js b/casRpsLogExtractor/casRpsLogExtractor.js index c1dbc53..74ceb0c 100644 --- a/casRpsLogExtractor/casRpsLogExtractor.js +++ b/casRpsLogExtractor/casRpsLogExtractor.js @@ -1,5 +1,5 @@ -//const sql = require("mssql/msnodesqlv8"); -const sql = require("mssql"); +const sql = require("mssql/msnodesqlv8"); +//const sql = require("mssql"); const request = require("request-promise"); const fs = require("fs"); const winston = require("winston"); @@ -40,10 +40,16 @@ const config = JSON.parse(fs.readFileSync("casRpsLogExtractor.config.json", "utf logger.log("debug", "Config" + JSON.stringify(config)) const processingType = "casrpslog"; +const maxResultsPerMessage=100; let maxPKey = ""; let maxPKeyUnsent = ""; -let lastOutMaxPKey = ""; let allRows = []; + +let currPollInterval = config.pollIntervalBusy; +let lastPollInterval; +let theInterval; + + try { lastOutMaxPKey = maxPKey = fs.readFileSync("casRpsLogExtractor.maxPKey.dat", "utf8"); logger.log("debug", "MaxPKey Loaded: " + maxPKey); @@ -62,9 +68,11 @@ sql.connect(config.sqlConfig, err => { const sqlRequest = new sql.Request(); sqlRequest.stream = true; - setInterval(_ => { - sqlRequest.query("select top 100 * from rpslog where pkey > '" + maxPKey + "' order by pkey"); - }, config.pollInterval) + function runNextPoll(){ + sqlRequest.query("select top " + maxResultsPerMessage + " * from rpslog where pkey > '" + maxPKey + "' order by pkey"); + } + + theInterval = setInterval(runNextPoll, currPollInterval); sqlRequest.on('row', row => { outData = { sourceTag: config.sourceTag, processingType: processingType, data: row }; @@ -90,6 +98,16 @@ sql.connect(config.sqlConfig, err => { logger.log("error","Sending Request to Service: " + err) }); process.stdout.write(">" + allRows.length + ">"); + if (allRows.length === maxResultsPerMessage) { + currPollInterval = config.pollIntervalBusy; + } else { + currPollInterval = config.pollIntervalIdle; + } + if (currPollInterval !== lastPollInterval) { + clearInterval(theInterval); + theInterval = setInterval(runNextPoll, currPollInterval); + } + allRows=[]; }