implemented "Gear shifting" in the casRpsLogExtractor.
This commit is contained in:
@@ -1,7 +1,8 @@
|
|||||||
{
|
{
|
||||||
"sourceTag": "CCERpsLog",
|
"sourceTag": "CCERpsLog",
|
||||||
"serverURI": "http://localhost:2222",
|
"serverURI": "http://localhost:2222",
|
||||||
"pollInterval": 10000,
|
"pollIntervalIdle": 10000,
|
||||||
|
"pollIntervalBusy": 1000,
|
||||||
"UP_sqlConfig": {
|
"UP_sqlConfig": {
|
||||||
"user": "casdbuser",
|
"user": "casdbuser",
|
||||||
"password": "4KcLao2016!",
|
"password": "4KcLao2016!",
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
//const sql = require("mssql/msnodesqlv8");
|
const sql = require("mssql/msnodesqlv8");
|
||||||
const sql = require("mssql");
|
//const sql = require("mssql");
|
||||||
const request = require("request-promise");
|
const request = require("request-promise");
|
||||||
const fs = require("fs");
|
const fs = require("fs");
|
||||||
const winston = require("winston");
|
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))
|
logger.log("debug", "Config" + JSON.stringify(config))
|
||||||
const processingType = "casrpslog";
|
const processingType = "casrpslog";
|
||||||
|
const maxResultsPerMessage=100;
|
||||||
let maxPKey = "";
|
let maxPKey = "";
|
||||||
let maxPKeyUnsent = "";
|
let maxPKeyUnsent = "";
|
||||||
let lastOutMaxPKey = "";
|
|
||||||
let allRows = [];
|
let allRows = [];
|
||||||
|
|
||||||
|
let currPollInterval = config.pollIntervalBusy;
|
||||||
|
let lastPollInterval;
|
||||||
|
let theInterval;
|
||||||
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
lastOutMaxPKey = maxPKey = fs.readFileSync("casRpsLogExtractor.maxPKey.dat", "utf8");
|
lastOutMaxPKey = maxPKey = fs.readFileSync("casRpsLogExtractor.maxPKey.dat", "utf8");
|
||||||
logger.log("debug", "MaxPKey Loaded: " + maxPKey);
|
logger.log("debug", "MaxPKey Loaded: " + maxPKey);
|
||||||
@@ -62,9 +68,11 @@ sql.connect(config.sqlConfig, err => {
|
|||||||
const sqlRequest = new sql.Request();
|
const sqlRequest = new sql.Request();
|
||||||
sqlRequest.stream = true;
|
sqlRequest.stream = true;
|
||||||
|
|
||||||
setInterval(_ => {
|
function runNextPoll(){
|
||||||
sqlRequest.query("select top 100 * from rpslog where pkey > '" + maxPKey + "' order by pkey");
|
sqlRequest.query("select top " + maxResultsPerMessage + " * from rpslog where pkey > '" + maxPKey + "' order by pkey");
|
||||||
}, config.pollInterval)
|
}
|
||||||
|
|
||||||
|
theInterval = setInterval(runNextPoll, currPollInterval);
|
||||||
|
|
||||||
sqlRequest.on('row', row => {
|
sqlRequest.on('row', row => {
|
||||||
outData = { sourceTag: config.sourceTag, processingType: processingType, data: 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)
|
logger.log("error","Sending Request to Service: " + err)
|
||||||
});
|
});
|
||||||
process.stdout.write(">" + allRows.length + ">");
|
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=[];
|
allRows=[];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user