Files
laermsensor-stack/readarchive/influx_post.js
T
admin 7c6692c040 readarchive: Import aus dem CSV-Archiv dazu
Gegenstueck zu readin/: waehrend readin alle 5 Minuten die Live-API
abfragt, liest readarchive rueckwirkend die Tagesdateien von
archive.sensor.community in dieselben Datenbanken ein.

Ohne History uebernommen (das Projekt lag bisher in einem eigenen
Repository unter Sensors/Laerm/laerm_readfromcsv_to_database).

Stand entspricht dort e33a7b7:
- Zeitstempel werden in UTC gespeichert, unabhaengig von der Zeitzone
  der Maschine. readin/parse.js macht das seit jeher richtig,
  readFromcsv.js lag zwei Stunden daneben.
- Fehlgeschlagene Influx-Writes werden gemeldet und setzen den
  Exit-Code, statt still verloren zu gehen.
- Der Container nimmt Parameter entgegen (ENTRYPOINT in Exec-Form).

Hinweis: mongo.js, influx_post.js und logit.js gibt es auch unter
readin/, mit abweichendem Stand. Zusammenfuehren waere der naechste
Schritt.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-31 06:36:01 +00:00

108 lines
3.3 KiB
JavaScript

// Access to influxDB vie HTTP
import axios from 'axios'
import { logit, logerror } from './logit.js'
import { DateTime } from 'luxon'
let INFLUXHOST = process.env.INFLUXHOST || "localhost"
let INFLUXPORT = process.env.INFLUXPORT || 8086
let INFLUXTOKEN = process.env.INFLUXTOKEN || 'empty'
let INFLUXDATABUCKET = process.env.INFLUXDATABUCKET || "sensor_data"
let INFLUXORG = process.env.INFLUXORG || "citysensor"
const INFLUXURL_READ = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/query?org=${INFLUXORG}`
const INFLUXURL_WRITE = `http://${INFLUXHOST}:${INFLUXPORT}/api/v2/write?org=${INFLUXORG}&bucket=${INFLUXDATABUCKET}&precision=ms`
// `${e}` liefert bei einem AggregateError nur "AggregateError" - die einzelnen
// Verbindungsfehler (je einer pro aufgelöster IP) stecken in e.errors bzw. e.cause.errors
function describeError(e) {
const sub = e.errors || e.cause?.errors
if (sub?.length) {
return `${e.message || e.name}: ` + sub.map(s => `${s.code} ${s.address}:${s.port}`).join(', ')
}
if (e.response) {
return `HTTP ${e.response.status} ${JSON.stringify(e.response.data)}`
}
return `${e.code ? e.code + ' ' : ''}${e.message || e}`
}
export const influxRead = async (query) => {
let start = DateTime.now()
let data = []
try {
let ret = await axios({
method: 'post',
url: INFLUXURL_READ,
data: query,
headers: {
Authorization: `Token ${INFLUXTOKEN}`,
Accept: 'application/csv',
'Content-type': 'application/vnd.flux'
},
timeout: 10000,
})
if (ret.status != 200) {
logerror(`doReadfromAPI Status: ${ret.status}`)
}
data = ret.data
} catch (e) {
logerror(`doReadfromAPI ${INFLUXURL_READ} ${describeError(e)}`)
}
logit(`ReadIn-Time: ${start.diffNow('seconds').toObject().seconds * -1} sec`)
return data
}
// liefert true, wenn Influx die Daten uebernommen hat, sonst false
export const influxWrite = async (data) => {
let start = DateTime.now()
let ok = false
try {
const ret = await axios({
method: 'post',
url: INFLUXURL_WRITE,
data: data,
headers: {
Authorization: `Token ${INFLUXTOKEN}`,
Accept: 'application/json',
'Content-Type': 'text/plain; charset=utf-8'
},
timeout: 10000,
})
if (ret.status != 204) {
logerror(`doWrite2API Status: ${ret.status}`)
} else {
ok = true
}
} catch (e) {
logerror(`doWrite2API ${INFLUXURL_WRITE} ${describeError(e)}`)
}
logit(`Influx-Write-Time: ${start.diffNow('seconds').toObject().seconds * -1} sec`)
return ok
}
/*
async function main() {
let data = `
pm,sid=140 P1=12,P2=13
pm,sid=142 P1=42,P2=13
pm,sid=143 P1=43,P2=13
pm,sid=144 P1=44,P2=13
thp,sid=141 temperature=23.5,humidity=48,pressure=998
`
let ret = await influxWrite(data)
process.exit()
let query = `from(bucket:"sensor_data")
|> range(start: -1mo)
|> filter(fn: (r) => r._measurement == "pm")
|> filter(fn: (r) => r.sid == "140")
`
let erg = await influxRead(query)
console.log(erg)
}
main().catch(console.error)
*/