From 887399ef858866f5e0d48ab977cd2d515d3f2da8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Reinhard=20X=2E=20F=C3=BCrst?= Date: Mon, 3 Aug 2026 10:39:27 +0000 Subject: [PATCH] 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 --- noisesensors/readarchive-stack.yml | 4 +- readarchive/package.json | 4 +- readarchive/readFromcsv.js | 69 ++++++++++++++++++++++-------- 3 files changed, 55 insertions(+), 22 deletions(-) diff --git a/noisesensors/readarchive-stack.yml b/noisesensors/readarchive-stack.yml index 6fad3df..1ef7104 100644 --- a/noisesensors/readarchive-stack.yml +++ b/noisesensors/readarchive-stack.yml @@ -16,7 +16,8 @@ # # readFromcsv.js liest jeden Parameter auch aus einer Env-Variablen: # -s = START (Default: gestern), -e = END (Default: Startdatum + 1 Tag, -# exklusiv), -t = TYP (Default: noise), -d = DBASE (Default: both). +# exklusiv), -t = TYP (Default: noise), -d = DBASE (Default: both), +# -p = PARALLEL (Default: 4) - so viele Sensoren werden gleichzeitig geholt. # # Voraussetzungen: # * Der Haupt-Stack laeuft, denn dessen Netz wird hier eingebunden. Der Name @@ -46,6 +47,7 @@ services: END: ${END:-} TYP: ${TYP:-noise} DBASE: ${DBASE:-influx} + PARALLEL: ${PARALLEL:-4} volumes: - ${LOCALDIR}/noisesensors/log:/var/log - ${LOCALDIR}/noisesensors/data/readarchive:/opt/app/readarchive/data diff --git a/readarchive/package.json b/readarchive/package.json index d4ca0a8..bec0f40 100644 --- a/readarchive/package.json +++ b/readarchive/package.json @@ -1,7 +1,7 @@ { "name": "sensors_readfromcsv", - "version": "3.1.2", - "date": "2023-12-23", + "version": "3.2.0", + "date": "2026-08-03", "description": "", "main": "readfromcsv.js", "scripts": { diff --git a/readarchive/readFromcsv.js b/readarchive/readFromcsv.js index cd82af3..b882521 100644 --- a/readarchive/readFromcsv.js +++ b/readarchive/readFromcsv.js @@ -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) {