Start using mongo timeseries
mooving all database relevant function to influx.js resp. mongo.js aktual data now also readable from mongo
This commit is contained in:
@@ -2,7 +2,6 @@
|
|||||||
import {DateTime} from "luxon"
|
import {DateTime} from "luxon"
|
||||||
import * as mongo from "../databases/mongo.js"
|
import * as mongo from "../databases/mongo.js"
|
||||||
import { returnOnError } from "../utilities/reporterror.js"
|
import { returnOnError } from "../utilities/reporterror.js"
|
||||||
import { fetchFromInflux } from "./getsensorData.js";
|
|
||||||
|
|
||||||
|
|
||||||
// Default distance for center search ( in km)
|
// Default distance for center search ( in km)
|
||||||
|
|||||||
@@ -136,25 +136,9 @@ export async function getSensorData(params) {
|
|||||||
|
|
||||||
// export const getActData = async (opts) => {
|
// export const getActData = async (opts) => {
|
||||||
export async function getActData(opts) {
|
export async function getActData(opts) {
|
||||||
let ret = {err: null, values: []}
|
let retI = await influx.fetchActData(opts)
|
||||||
let sorting = ''
|
let retM = await mongo.fetchActData(opts)
|
||||||
if(opts.sort) {
|
return retI
|
||||||
if (opts.sort === 1) {
|
|
||||||
sorting = '|> sort(columns: ["_time"], desc: false)'
|
|
||||||
} else if (opts.sort === -1) {
|
|
||||||
sorting = '|> sort(columns: ["_time"], desc: true)'
|
|
||||||
}
|
|
||||||
}
|
|
||||||
// build the flux query
|
|
||||||
let query = `
|
|
||||||
from(bucket: "sensor_data")
|
|
||||||
|> range(${opts.start}, ${opts.stop})
|
|
||||||
|> filter(fn: (r) => r.sid == "${opts.sensorid}")
|
|
||||||
${sorting}
|
|
||||||
|> keep(columns: ["_time","_field","_value"])
|
|
||||||
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
|
||||||
`
|
|
||||||
return await fetchFromInflux(ret, query)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -234,21 +218,7 @@ export var getLongAvg = async (params) => {
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
export const fetchFromInflux = async (ret, query) => {
|
|
||||||
let { values, err} = await influx.influxRead(query)
|
|
||||||
if(err) {
|
|
||||||
if(err.toString().includes('400')) {
|
|
||||||
return returnOnError(ret, 'SYNTAXURL', fetchFromInflux.name)
|
|
||||||
} else {
|
|
||||||
return returnOnError(ret, err, fetchFromInflux.name)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (values.length <= 2) {
|
|
||||||
return returnOnError(ret, 'NODATA', fetchFromInflux.name)
|
|
||||||
}
|
|
||||||
ret.values = csv2Json(values)
|
|
||||||
return ret
|
|
||||||
}
|
|
||||||
|
|
||||||
// *********************************************
|
// *********************************************
|
||||||
// function getMAPaktData() {
|
// function getMAPaktData() {
|
||||||
|
|||||||
+97
-2
@@ -5,6 +5,7 @@ import { DateTime } from 'luxon'
|
|||||||
// import csvParse from 'csv-parser'
|
// import csvParse from 'csv-parser'
|
||||||
import { logit, logerror } from '../utilities/logit.js'
|
import { logit, logerror } from '../utilities/logit.js'
|
||||||
import {returnOnError} from "../utilities/reporterror.js";
|
import {returnOnError} from "../utilities/reporterror.js";
|
||||||
|
import {csv2Json} from "../utilities/csv2json.js";
|
||||||
|
|
||||||
let INFLUXHOST = process.env.INFLUXHOST || "localhost"
|
let INFLUXHOST = process.env.INFLUXHOST || "localhost"
|
||||||
let INFLUXPORT = process.env.INFLUXPORT || 8086
|
let INFLUXPORT = process.env.INFLUXPORT || 8086
|
||||||
@@ -17,7 +18,7 @@ let INFLUXORG = process.env.INFLUXORG || "citysensor"
|
|||||||
const INFLUXURL_READ = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/query?org=${INFLUXORG}`
|
const INFLUXURL_READ = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/query?org=${INFLUXORG}`
|
||||||
const INFLUXURL_WRITE = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/write?org=${INFLUXORG}&bucket=${INFLUXDATABUCKET}`
|
const INFLUXURL_WRITE = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/write?org=${INFLUXORG}&bucket=${INFLUXDATABUCKET}`
|
||||||
|
|
||||||
export const influxRead = async (query) => {
|
const influxRead = async (query) => {
|
||||||
let start = DateTime.now()
|
let start = DateTime.now()
|
||||||
logit(`ReadInflux from ${INFLUXURL_READ}`)
|
logit(`ReadInflux from ${INFLUXURL_READ}`)
|
||||||
let erg = { values: [], err: null}
|
let erg = { values: [], err: null}
|
||||||
@@ -45,7 +46,7 @@ export const influxRead = async (query) => {
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
export const influxWrite = async (data) => {
|
const influxWrite = async (data) => {
|
||||||
let start = DateTime.now()
|
let start = DateTime.now()
|
||||||
let ret
|
let ret
|
||||||
try {
|
try {
|
||||||
@@ -69,3 +70,97 @@ export const influxWrite = async (data) => {
|
|||||||
logit(`Influx-Write-Time: ${start.diffNow('seconds').toObject().seconds * -1} sec`)
|
logit(`Influx-Write-Time: ${start.diffNow('seconds').toObject().seconds * -1} sec`)
|
||||||
return ret
|
return ret
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const fetchFromInflux = async (ret, query) => {
|
||||||
|
let { values, err} = await influxRead(query)
|
||||||
|
if(err) {
|
||||||
|
if(err.toString().includes('400')) {
|
||||||
|
return returnOnError(ret, 'SYNTAXURL', fetchFromInflux.name)
|
||||||
|
} else {
|
||||||
|
return returnOnError(ret, err, fetchFromInflux.name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (values.length <= 2) {
|
||||||
|
return returnOnError(ret, 'NODATA', fetchFromInflux.name)
|
||||||
|
}
|
||||||
|
ret.values = csv2Json(values)
|
||||||
|
return ret
|
||||||
|
}
|
||||||
|
|
||||||
|
export const fetchActData = async (opts) => {
|
||||||
|
let ret = {err: null, values: []}
|
||||||
|
let sorting = ''
|
||||||
|
if(opts.sort) {
|
||||||
|
if (opts.sort === 1) {
|
||||||
|
sorting = '|> sort(columns: ["_time"], desc: false)'
|
||||||
|
} else if (opts.sort === -1) {
|
||||||
|
sorting = '|> sort(columns: ["_time"], desc: true)'
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// build the flux query
|
||||||
|
let query = `
|
||||||
|
from(bucket: "sensor_data")
|
||||||
|
|> range(${opts.start}, ${opts.stop})
|
||||||
|
|> filter(fn: (r) => r.sid == "${opts.sensorid}")
|
||||||
|
${sorting}
|
||||||
|
|> keep(columns: ["_time","_field","_value"])
|
||||||
|
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||||
|
`
|
||||||
|
return await fetchFromInflux(ret, query)
|
||||||
|
}
|
||||||
|
|
||||||
|
export const fetchNoiseAVGData = async (opts) => {
|
||||||
|
let ret = {err: null, values: []}
|
||||||
|
let small = '|> keep(columns: ["_time", "peakcount", "n_AVG"])'
|
||||||
|
if (opts.long) {
|
||||||
|
small = ''
|
||||||
|
}
|
||||||
|
let queryAVG = `
|
||||||
|
import "math"
|
||||||
|
threshold = ${opts.peak}
|
||||||
|
|
||||||
|
data = from(bucket: "sensor_data")
|
||||||
|
|> range(${opts.start}, ${opts.stop})
|
||||||
|
|> filter(fn: (r) => r["sid"] == "${opts.sensorid}")
|
||||||
|
e10 = data
|
||||||
|
|> filter(fn: (r) => r._field == "E10tel_eq")
|
||||||
|
|> aggregateWindow(every: 1h, fn: mean, createEmpty: false)
|
||||||
|
|> map(fn: (r) => ({r with _value: (10.0 * math.log10(x: r._value))}))
|
||||||
|
|> keep(columns: ["_time","_field","_value"])
|
||||||
|
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||||
|
|> rename(columns: {"E10tel_eq" : "n_AVG"})
|
||||||
|
ecnt = data
|
||||||
|
|> filter(fn: (r) => r._field == "E10tel_eq")
|
||||||
|
|> aggregateWindow(every: 1h, fn: count, createEmpty: false)
|
||||||
|
|> keep(columns: ["_time","_field","_value"])
|
||||||
|
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||||
|
|> rename(columns: {"E10tel_eq" : "count"})
|
||||||
|
esum = data
|
||||||
|
|> filter(fn: (r) => r._field == "E10tel_eq")
|
||||||
|
|> aggregateWindow(every: 1h, fn: sum, createEmpty: false)
|
||||||
|
|> keep(columns: ["_time","_field","_value"])
|
||||||
|
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||||
|
|> rename(columns: {"E10tel_eq" : "n_sum"})
|
||||||
|
peak = data
|
||||||
|
|> filter(fn: (r) => r._field == "noise_LA_max")
|
||||||
|
|> aggregateWindow(
|
||||||
|
every: 1h,
|
||||||
|
fn: (column, tables=<-) => tables
|
||||||
|
|> reduce(
|
||||||
|
identity: {peakcount: 0.0},
|
||||||
|
fn: (r, accumulator) => ({
|
||||||
|
peakcount: if r._value >= threshold then
|
||||||
|
accumulator.peakcount + 1.0
|
||||||
|
else
|
||||||
|
accumulator.peakcount + 0.0,
|
||||||
|
}),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|> keep(columns: ["_time","peakcount"])
|
||||||
|
part1 = join( tables: {e10: e10, ecnt: ecnt}, on: ["_time"])
|
||||||
|
part2 = join( tables: {esum: esum, peak: peak}, on: ["_time"])
|
||||||
|
join( tables: {P1: part1, P2: part2}, on: ["_time"])
|
||||||
|
${small}
|
||||||
|
`
|
||||||
|
return await fetchFromInflux(ret, queryAVG)
|
||||||
|
}
|
||||||
+55
-1
@@ -30,7 +30,7 @@ export const connectMongo = async () => {
|
|||||||
// logit(`Try to connect to ${MONGO_URL}`)
|
// logit(`Try to connect to ${MONGO_URL}`)
|
||||||
// let client = await MongoClient.connect(MONGO_URL, { useNewUrlParser: true , useUnifiedTopology: true })
|
// let client = await MongoClient.connect(MONGO_URL, { useNewUrlParser: true , useUnifiedTopology: true })
|
||||||
let client = await MongoClient.connect(MONGO_URL)
|
let client = await MongoClient.connect(MONGO_URL)
|
||||||
// logit(`Mongodbase connected to ${MONGO_URL}`)
|
logit(`Mongodbase connected to ${MONGO_URL}`)
|
||||||
return client
|
return client
|
||||||
}
|
}
|
||||||
catch(error){
|
catch(error){
|
||||||
@@ -143,3 +143,57 @@ export const readAKWs = async (options) => {
|
|||||||
}
|
}
|
||||||
return ret
|
return ret
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export const fetchActData = async (opts) => {
|
||||||
|
let ret = {err: null, values: []}
|
||||||
|
let start = opts.start.slice(7)
|
||||||
|
let end = opts.stop.slice(6)
|
||||||
|
start = DateTime.fromISO(start).toJSDate()
|
||||||
|
end = DateTime.fromISO(end).toJSDate()
|
||||||
|
let query = {sensorid: opts.sensorid, datetime: {$gte: start, $lt: end}}
|
||||||
|
let client = await connectMongo()
|
||||||
|
try {
|
||||||
|
ret.values = await client.db(MONGOBASE).collection('sensors')
|
||||||
|
.find(query).toArray()
|
||||||
|
}
|
||||||
|
catch(e) {
|
||||||
|
ret.err = e
|
||||||
|
}
|
||||||
|
finally {
|
||||||
|
client.close()
|
||||||
|
}
|
||||||
|
return ret
|
||||||
|
}
|
||||||
|
/*
|
||||||
|
let docs = await collection.find(
|
||||||
|
{ datetime:
|
||||||
|
{ $gte: start.toDate(), $lt: end.toDate() }
|
||||||
|
},
|
||||||
|
{ projection:
|
||||||
|
{_id:0, E_eq:0, E_mx:0, E_mi:0, E10tel_mx:0, E10tel_mi:0}, sort: {datetime: sort}
|
||||||
|
},
|
||||||
|
).toArray();
|
||||||
|
*/
|
||||||
|
|
||||||
|
|
||||||
|
export const fetchActDataxx = async (opts) => {
|
||||||
|
let ret = {err: null, values: []}
|
||||||
|
let sorting = ''
|
||||||
|
if(opts.sort) {
|
||||||
|
if (opts.sort === 1) {
|
||||||
|
sorting = '|> sort(columns: ["_time"], desc: false)'
|
||||||
|
} else if (opts.sort === -1) {
|
||||||
|
sorting = '|> sort(columns: ["_time"], desc: true)'
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// build the flux query
|
||||||
|
let query = `
|
||||||
|
from(bucket: "sensor_data")
|
||||||
|
|> range(${opts.start}, ${opts.stop})
|
||||||
|
|> filter(fn: (r) => r.sid == "${opts.sensorid}")
|
||||||
|
${sorting}
|
||||||
|
|> keep(columns: ["_time","_field","_value"])
|
||||||
|
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
||||||
|
`
|
||||||
|
return await fetchFromInflux(ret, query)
|
||||||
|
}
|
||||||
|
|||||||
+8
-55
@@ -2,10 +2,12 @@
|
|||||||
// rxf 2023-03-05
|
// rxf 2023-03-05
|
||||||
|
|
||||||
import {returnOnError} from "../utilities/reporterror.js";
|
import {returnOnError} from "../utilities/reporterror.js";
|
||||||
import { getActData, getAvgData, getLongAvg, fetchFromInflux, calcRange} from "../actions/getsensorData.js"
|
import { getActData, getAvgData, getLongAvg, calcRange} from "../actions/getsensorData.js"
|
||||||
import checkParams from "../utilities/checkparams.js";
|
import checkParams from "../utilities/checkparams.js";
|
||||||
import {DateTime} from 'luxon'
|
import {DateTime} from 'luxon'
|
||||||
import { translate as trans } from '../routes/api.js'
|
import { translate as trans } from '../routes/api.js'
|
||||||
|
import * as influx from "../databases/influx.js"
|
||||||
|
import * as mongo from "../databases/mongo.js"
|
||||||
|
|
||||||
|
|
||||||
const setoptionfromtable = (opt,tabval) => {
|
const setoptionfromtable = (opt,tabval) => {
|
||||||
@@ -429,61 +431,8 @@ const getAPIprops = (opt) => {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const getNoiseAVGData = async (opts) => {
|
const getNoiseAVGData = async (opts) => {
|
||||||
let ret = {err: null, values: []}
|
let ret = await influx.fetchNoiseAVGData(opts)
|
||||||
let emptyValues = {n_AVG: -1}
|
|
||||||
let small = '|> keep(columns: ["_time", "peakcount", "n_AVG"])'
|
|
||||||
if (opts.long) {
|
|
||||||
small = ''
|
|
||||||
emptyValues.n_sum = -1
|
|
||||||
}
|
|
||||||
let queryAVG = `
|
|
||||||
import "math"
|
|
||||||
threshold = ${opts.peak}
|
|
||||||
|
|
||||||
data = from(bucket: "sensor_data")
|
|
||||||
|> range(${opts.start}, ${opts.stop})
|
|
||||||
|> filter(fn: (r) => r["sid"] == "${opts.sensorid}")
|
|
||||||
e10 = data
|
|
||||||
|> filter(fn: (r) => r._field == "E10tel_eq")
|
|
||||||
|> aggregateWindow(every: 1h, fn: mean, createEmpty: false)
|
|
||||||
|> map(fn: (r) => ({r with _value: (10.0 * math.log10(x: r._value))}))
|
|
||||||
|> keep(columns: ["_time","_field","_value"])
|
|
||||||
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
|
||||||
|> rename(columns: {"E10tel_eq" : "n_AVG"})
|
|
||||||
ecnt = data
|
|
||||||
|> filter(fn: (r) => r._field == "E10tel_eq")
|
|
||||||
|> aggregateWindow(every: 1h, fn: count, createEmpty: false)
|
|
||||||
|> keep(columns: ["_time","_field","_value"])
|
|
||||||
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
|
||||||
|> rename(columns: {"E10tel_eq" : "count"})
|
|
||||||
esum = data
|
|
||||||
|> filter(fn: (r) => r._field == "E10tel_eq")
|
|
||||||
|> aggregateWindow(every: 1h, fn: sum, createEmpty: false)
|
|
||||||
|> keep(columns: ["_time","_field","_value"])
|
|
||||||
|> pivot(rowKey:["_time"], columnKey: ["_field"], valueColumn: "_value")
|
|
||||||
|> rename(columns: {"E10tel_eq" : "n_sum"})
|
|
||||||
peak = data
|
|
||||||
|> filter(fn: (r) => r._field == "noise_LA_max")
|
|
||||||
|> aggregateWindow(
|
|
||||||
every: 1h,
|
|
||||||
fn: (column, tables=<-) => tables
|
|
||||||
|> reduce(
|
|
||||||
identity: {peakcount: 0.0},
|
|
||||||
fn: (r, accumulator) => ({
|
|
||||||
peakcount: if r._value >= threshold then
|
|
||||||
accumulator.peakcount + 1.0
|
|
||||||
else
|
|
||||||
accumulator.peakcount + 0.0,
|
|
||||||
}),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
|> keep(columns: ["_time","peakcount"])
|
|
||||||
part1 = join( tables: {e10: e10, ecnt: ecnt}, on: ["_time"])
|
|
||||||
part2 = join( tables: {esum: esum, peak: peak}, on: ["_time"])
|
|
||||||
join( tables: {P1: part1, P2: part2}, on: ["_time"])
|
|
||||||
${small}
|
|
||||||
`
|
|
||||||
ret = await fetchFromInflux(ret, queryAVG)
|
|
||||||
if(ret.err) {
|
if(ret.err) {
|
||||||
return returnOnError(ret, ret.err, getNoiseAVGData.name)
|
return returnOnError(ret, ret.err, getNoiseAVGData.name)
|
||||||
}
|
}
|
||||||
@@ -496,6 +445,10 @@ peak = data
|
|||||||
// hour in an element in docs becomes the index into the new array (for every new day this
|
// hour in an element in docs becomes the index into the new array (for every new day this
|
||||||
// index will be incremented by 24). Missing values are marked by: {n_sum=-1, n_AVG=-1}.
|
// index will be incremented by 24). Missing values are marked by: {n_sum=-1, n_AVG=-1}.
|
||||||
// For havg add the missed hours to the arry
|
// For havg add the missed hours to the arry
|
||||||
|
let emptyValues = {n_AVG: -1}
|
||||||
|
if (opts.long) {
|
||||||
|
emptyValues.n_sum = -1
|
||||||
|
}
|
||||||
const misshours = DateTime.fromISO(ret.values[0].datetime).get('hour')
|
const misshours = DateTime.fromISO(ret.values[0].datetime).get('hour')
|
||||||
let hoursArr = new Array(opts.span * 24 + misshours); // generate new array
|
let hoursArr = new Array(opts.span * 24 + misshours); // generate new array
|
||||||
hoursArr.fill(emptyValues) // fill array with 'empty' values
|
hoursArr.fill(emptyValues) // fill array with 'empty' values
|
||||||
|
|||||||
Reference in New Issue
Block a user