readarchive: Sensoren eines Tages parallel abholen

Die Laufzeit steckt fast vollstaendig im Warten auf die CSV-Dateien von
archive.sensor.community - je Sensor rund 0,24 s, davon nur 0,02 s der
Influx-Write und 0,01 s die Mongo-Abfrage. Die Schleife ueber die Sensoren
laeuft deshalb jetzt mit mehreren Arbeitern gleichzeitig (-p bzw. PARALLEL,
Default 4, 1 = altes Verhalten).

Gemessen fuer 2026-08-01 (280 Laermsensoren, 344660 Werte, -d influx):
1 -> 79 s, 4 -> 33 s, 8 -> 35 s, 16 -> 31 s. Ab etwa 4 gleichzeitigen
Abrufen liefert das Archiv nicht mehr schneller; derselbe Verlauf zeigt
sich mit blossem curl ohne Datenbank (71 / 32 / 28 / 29 / 29 s bei 1, 4,
8, 16, 32 Abrufen). Der Default steht deshalb auf 4 und nicht hoeher -
mehr erzeugt nur Last beim fremden Server.

checkProperties() wird jetzt awaited - der Aufruf lief bisher ohne await
gegen das bulkWrite am Ende von main(), was mit parallelen Sensoren kein
Randfall mehr waere.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-03 10:39:27 +00:00
parent eeab1ddbe5
commit 887399ef85
3 changed files with 55 additions and 22 deletions
+50 -19
View File
@@ -6,6 +6,9 @@
// Version:
//
// V 3.2.0 2026-08-03 rxf
// - Sensoren eines Tages werden parallel abgearbeitet (-p / PARALLEL)
//
// V 3.1.1 2023-11-14 rxf
// - Enddatum eingeführt
//
@@ -23,7 +26,12 @@
const TYP = process.env.TYP || 'noise'
const START = process.env.START
const END = process.env.END
const DBASE = process.env.DBASE || 'both'
const DBASE = process.env.DBASE || 'both'
// Anzahl der Sensoren, die gleichzeitig geholt werden. Die Laufzeit steckt fast
// vollstaendig im Warten auf archive.sensor.community (~0,15 s je CSV), nicht in
// Mongo/Influx - deshalb bringt Parallelitaet hier den Hebel. 1 = altes Verhalten.
// Mehr als 4 bringt nichts: ab da liefert das Archiv nicht mehr schneller.
const PARALLEL = parseInt(process.env.PARALLEL) || 4
const DATABASE = ['mongo', 'influx', 'both']
@@ -228,29 +236,42 @@ async function readSensorsperDay(client, args) {
return insertedCount
}
// 'limit' Arbeiter teilen sich die Liste: jeder nimmt sich den naechsten freien
// Eintrag, sobald er fertig ist. Kein Batching, damit ein langsamer Abruf die
// anderen nicht aufhaelt. Node ist single-threaded - die gemeinsamen Zaehler und
// Arrays brauchen deshalb keine Sperre, es laeuft immer nur ein Stueck Code
// zwischen zwei await-Punkten.
const runPool = async (items, limit, worker) => {
let next = 0
const runner = async () => {
while (next < items.length) {
const i = next++
await worker(items[i])
}
}
await Promise.all(Array.from({length: Math.min(limit, items.length)}, runner))
}
const scanOneSensor = async (client, sensorList, dt, args) => {
let inserted = 0
let missed = []
let listcnt = sensorList.length
for (let item of sensorList) { // iterate the list
let list = sensorList
if (args.sensornbr !== undefined) {
list = list.filter((item) => item.sensorid === args.sensornbr)
}
if (args.nbrOfEntries !== undefined) {
list = list.slice(0, args.nbrOfEntries)
args.nbrOfEntries -= list.length
}
let listcnt = list.length
await runPool(list, args.parallel, async (item) => {
let sid = item.sensorid
if(args.sensornbr !== undefined) {
if(args.sensornbr !== sid) {
continue
}
}
logit(`${item.url}`)
let icount = await putOneSensorInDb(client, item.url, sid, dt, missed, item.indoor, item.sensorType, args) // put one sensor data into DB
logit(`${listcnt} -- Sensor ${sid}: ${icount} Einträge`)
inserted += icount
listcnt--
if (args.nbrOfEntries !== undefined) {
if (--args.nbrOfEntries === 0) {
break
}
}
}
})
return {inserted: inserted, missed: missed}
}
@@ -343,7 +364,9 @@ async function readOneSensorOneDay(client, name, sid, dt, missed, indoor, args)
dataline += `${typ},sid=${entry.sensorid} ${influxvalue} ${entry.datetime.valueOf()}\n`
}
}
checkProperties(client, firstItem, indoor, sname, typ, dt, newProps)
// await, weil newProps sonst mit dem bulkWrite am Ende von main() um die
// Wette laeuft - bei parallelen Sensoren ist das kein Randfall mehr
await checkProperties(client, firstItem, indoor, sname, typ, dt, newProps)
} catch (e) {
console.log(e)
missed.push(name)
@@ -422,10 +445,10 @@ async function enterOneSensorinDB(client, sid, erg, dataline, day, styp, args) {
// Parse command line options
function parse_cmdline(argv) {
let parser = new mod_getopt.BasicParser('s:(start)e:(end)t:(type)h(help)v(version)a:(entries)n:(sensorid)d:(dbase)',argv);
let parser = new mod_getopt.BasicParser('s:(start)e:(end)t:(type)h(help)v(version)a:(entries)n:(sensorid)d:(dbase)p:(parallel)',argv);
let option;
let std = START ? START : DateTime.now().startOf('day').minus({days: 1}).toFormat('yyyy-LL-dd'); // yesterday
let ret = {startDate: std, endDate: '', sensorType: TYP, database: DBASE};
let ret = {startDate: std, endDate: '', sensorType: TYP, database: DBASE, parallel: PARALLEL};
while((option = parser.getopt()) !== undefined) {
switch(option.option) {
@@ -460,6 +483,13 @@ function parse_cmdline(argv) {
}
break;
case 'p':
let par = parseInt(option.optarg.trim())
if (par > 0) {
ret.parallel = par
}
break;
case 'v':
console.log(`Version: ${pkg.version} from ${pkg.date}`);
console.log();
@@ -467,7 +497,7 @@ function parse_cmdline(argv) {
break;
case 'h':
console.log("Usage: node readFromcvs.js [-h] [-s startDate] [-e endDate] [-t sensorType] [-v version] [-h help]")
console.log("Usage: node readFromcvs.js [-h] [-s startDate] [-e endDate] [-t sensorType] [-d database] [-p parallel] [-v version] [-h help]")
if (DEVELOP) {
console.log(" [-a entries] [-n sensorid]");
}
@@ -478,6 +508,7 @@ function parse_cmdline(argv) {
console.log(" -t sensorType: if given, only those sensors will be used; default: 'noise'");
console.log(" allowed types: 'noise', 'radiactivity', 'pm' and 'thp'")
console.log(" -d database: 'mongo', 'influx' or 'both'; default: 'both'")
console.log(` -p parallel: number of sensors fetched at the same time; default: ${PARALLEL}`)
console.log(" -v version: show version");
console.log(" -h this help text");
if(DEVELOP) {