From 7a665f645b594352068c8be7228157d213273c32 Mon Sep 17 00:00:00 2001 From: Antoine Pultier Date: Fri, 2 Oct 2026 22:46:12 +0200 Subject: [PATCH 1/4] =?UTF-8?q?docs:=20=F0=9F=93=9D=20close=20the=20ClickH?= =?UTF-8?q?ouse=20pre-production=20readiness=20task?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PR 44 is merged and the CI is green on main. What is left is staging, the operations drills and the release, which the release plan tracks. Co-Authored-By: Claude Sonnet 5.5 --- .../clickhouse-preproduction-readiness.md | 7 +++++++ 1 file changed, 7 insertions(+) rename {current_tasks => done}/clickhouse-preproduction-readiness.md (94%) diff --git a/current_tasks/clickhouse-preproduction-readiness.md b/done/clickhouse-preproduction-readiness.md similarity index 94% rename from current_tasks/clickhouse-preproduction-readiness.md rename to done/clickhouse-preproduction-readiness.md index 73f64115..0ca4b9f2 100644 --- a/current_tasks/clickhouse-preproduction-readiness.md +++ b/done/clickhouse-preproduction-readiness.md @@ -82,6 +82,13 @@ Databases created by earlier builds are refused at startup with a clear message No known data-correctness bug, failures reported with the right status and recovering without a restart, a documented backup and restore that was practised, documentation of what is and is not guaranteed, and evidence from a real service instead of mocks. That is met for the code. What remains of the release plan is evidence from a real environment: staging the image and the chart, exercising operations there, and publishing (steps 2 to 4 of `docs/PREPRODUCTION_RELEASE_PLAN.md`), plus a green CI run. +## Closed (2 Oct 2026) + +Merged to `main` with PR 44, and the CI/CD Pipeline is green on the merge commit (backend matrix, frontend, +Python SDK, audit, Helm, Docker smoke, live Prometheus compatibility): step 1 of the release plan has its +evidence. What is left is not code: staging, operations drills and the release itself (steps 2 to 4 of +`docs/PREPRODUCTION_RELEASE_PLAN.md`), plus the frontend generator upgrade (step 5). + ## Notes - Keep tests generic where practical, but accept backend-specific validation where operational behavior differs. From 63524fab8925a3c5950e920efa1b1abe20c30c3a Mon Sep 17 00:00:00 2001 From: Antoine Pultier Date: Fri, 2 Oct 2026 22:46:13 +0200 Subject: [PATCH 2/4] =?UTF-8?q?docs:=20=F0=9F=93=9D=20plan=20the=20DuckDB?= =?UTF-8?q?=20bulk=20writes,=20with=20the=20baseline=20numbers?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Sonnet 5.5 --- current_tasks/duckdb-bulk-writes.md | 46 +++++++++++++++++++++++++++++ ideas/duckdb-bulk-writes.md | 16 ---------- 2 files changed, 46 insertions(+), 16 deletions(-) create mode 100644 current_tasks/duckdb-bulk-writes.md delete mode 100644 ideas/duckdb-bulk-writes.md diff --git a/current_tasks/duckdb-bulk-writes.md b/current_tasks/duckdb-bulk-writes.md new file mode 100644 index 00000000..480a9e1b --- /dev/null +++ b/current_tasks/duckdb-bulk-writes.md @@ -0,0 +1,46 @@ +# DuckDB: bulk writes + +## Goal + +Writing many series, or many strings, to DuckDB must cost a handful of statements whatever the number of +series, like PostgreSQL, TimescaleDB and ClickHouse. DuckDB is not the main target (local analysis), but the +per-series path is slow and the `#[cached]` ids have the rollback problem the other backends lost. + +## Baseline (2 Oct 2026, release build, macOS, `tests/perf/scale.sh 3000` and `tests/perf/strings.sh 3000`) + +| write | before | +|---|---| +| 3000 new series (30 000 samples) | 2.69 s | +| the same series again | 0.37 s | +| strings: 3000 series x 10, 50 distinct strings | 2.51 s | +| strings: one series of 10 000 distinct strings | 1.50 s | +| strings: the first request again | 0.37 s | + +Why: `get_sensor_id_or_create_sensor` is one `SELECT`, one `INSERT` and one `INSERT` per label (plus the +dictionaries) per sensor, the string dictionary is one statement per distinct string, and every sensor opens +its own appender on its value table. All of it behind `#[cached]` wrappers that are filled inside the +transaction, so a rollback leaves ids of sensors that were never committed (until the 120 s TTL ends). + +## Plan + +1. `duckdb_registration.rs`: register all the sensors of a batch with a few statements (ids by chunked + `IN`, appenders for units, label dictionaries, sensors and labels). Nothing cached. +2. Bulk string dictionary for the whole batch. +3. One appender per value table for the whole batch instead of one per sensor. +4. Remove the `#[cached]` wrappers, `forget_sensor_id` and their cache clears. +5. Before/after numbers, tests on top of the backend-generic ones, docs. + +## Done when + +- [ ] Before/after numbers below. +- [ ] A failed batch leaves no sensor, label or string behind and a later write of the same series works. +- [ ] Backend-generic tests still pass on DuckDB, and on SQLite and TimescaleDB where they are generic. +- [ ] Full suites, clippy, and the DuckDB suite. + +## Not in scope + +The aggregated selector read (`BulkSelectorBackend::read_aggregated_samples`) is still one query per series on +DuckDB (0.30 s for a 100 series remote read with a step). The `time_bucket` SQL of `duckdb_bucketed_cte` +extends to many sensors with `GROUP BY sensor_id, bucket`. A read, not a write: separate task if wanted. + +## Progress diff --git a/ideas/duckdb-bulk-writes.md b/ideas/duckdb-bulk-writes.md deleted file mode 100644 index fae1eb62..00000000 --- a/ideas/duckdb-bulk-writes.md +++ /dev/null @@ -1,16 +0,0 @@ -# DuckDB: register many new series in bulk - -Measured with `tests/perf/scale.sh 3000` on a release build (after the bulk read work): - -- write of 3000 new series (30 000 samples): 2.6 to 3.0 s; the same series again: 0.35 to 0.39 s -- strings, 3000 series x 10 samples, 50 distinct strings: 4.5 s; one series of 10 000 distinct strings: 2.1 s - -Series registration (`get_sensor_id_or_create_sensor`) is one `SELECT`, one `INSERT` and one `INSERT` per label, -and the string dictionary is one statement per distinct string. The value tables already use appenders. The -same pattern as `pg_sensor_registration` (sort, multi-row insert, read the ids back) would apply with -appenders on `sensors`, `labels` and the dictionaries. - -Not done because DuckDB is the "less mature" backend (local analysis, not ingestion); worth doing if DuckDB -becomes an ingestion target. The aggregated selector read (`BulkSelectorBackend::read_aggregated_samples`) -is also still one query per series on DuckDB (0.3 s for a 100 series remote read with a step); the -`time_bucket` SQL of `duckdb_bucketed_cte` extends to many sensors with `GROUP BY sensor_id, bucket`. From 8c9e9402bb516d55ffa9106177dcf81eb8d905ee Mon Sep 17 00:00:00 2001 From: Antoine Pultier Date: Fri, 2 Oct 2026 22:52:09 +0200 Subject: [PATCH 3/4] =?UTF-8?q?perf:=20=E2=9A=A1=EF=B8=8F=20write=20batche?= =?UTF-8?q?s=20to=20DuckDB=20with=20a=20handful=20of=20statements?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Register the sensors, units, labels and strings of a whole batch with chunked lookups and appenders instead of a few statements per sensor and per string, and write the samples through one appender per value table instead of one per sensor. 3000 new series: 2.69 s to 0.23 s; 3000 series of strings: 2.51 s to 0.20 s. Nothing is cached any more, so a rolled back batch leaves no stale sensor id, and the cache clearing of deletions and test cleanup is gone. Co-Authored-By: Claude Sonnet 5.5 --- src/storage/duckdb/duckdb_publishers.rs | 150 ++++++--- src/storage/duckdb/duckdb_registration.rs | 388 ++++++++++++++++++++++ src/storage/duckdb/duckdb_utilities.rs | 167 ---------- src/storage/duckdb/mod.rs | 68 +--- 4 files changed, 505 insertions(+), 268 deletions(-) create mode 100644 src/storage/duckdb/duckdb_registration.rs delete mode 100644 src/storage/duckdb/duckdb_utilities.rs diff --git a/src/storage/duckdb/duckdb_publishers.rs b/src/storage/duckdb/duckdb_publishers.rs index 77f84a32..aee2b3c8 100644 --- a/src/storage/duckdb/duckdb_publishers.rs +++ b/src/storage/duckdb/duckdb_publishers.rs @@ -1,124 +1,200 @@ -use super::duckdb_utilities::get_string_value_id_or_create; -use crate::datamodel::Sample; -use anyhow::Result; -use duckdb::{Transaction, params}; +use super::duckdb_registration::{ensure_string_ids, register_sensors}; +use crate::datamodel::batch::SingleSensorBatch; +use crate::datamodel::{Sample, TypedSamples}; +use anyhow::{Context, Result}; +use duckdb::{Appender, Connection, params}; use geo::Point; use rust_decimal::Decimal; use serde_json::Value; +use std::collections::{BTreeSet, HashMap}; -pub fn publish_integer_values( - transaction: &Transaction, +/// Writes a batch: the sensors and the strings are registered with a handful of statements, then +/// the samples of every sensor go through one appender per value table. +pub fn publish_batch(connection: &Connection, sensors: &[SingleSensorBatch]) -> Result<()> { + let sensor_refs: Vec<&crate::datamodel::Sensor> = + sensors.iter().map(|batch| batch.sensor.as_ref()).collect(); + let ids = register_sensors(connection, &sensor_refs)?; + + let guards: Vec<_> = sensors + .iter() + .map(|batch| batch.samples.blocking_read()) + .collect(); + let mut strings: BTreeSet<&str> = BTreeSet::new(); + for guard in &guards { + if let TypedSamples::String(samples) = &**guard { + strings.extend(samples.iter().map(|sample| sample.value.as_str())); + } + } + let string_ids = ensure_string_ids(connection, &strings)?; + + let mut appenders = Appenders::new(connection); + for (batch, guard) in sensors.iter().zip(&guards) { + let sensor_id = *ids + .get(&batch.sensor.uuid) + .context("a registered sensor has no id")?; + match &**guard { + TypedSamples::Integer(samples) => { + publish_integer_values(appenders.get("integer_values")?, sensor_id, samples)? + } + TypedSamples::Numeric(samples) => { + publish_numeric_values(appenders.get("numeric_values")?, sensor_id, samples)? + } + TypedSamples::Float(samples) => { + publish_float_values(appenders.get("float_values")?, sensor_id, samples)? + } + TypedSamples::String(samples) => publish_string_values( + appenders.get("string_values")?, + sensor_id, + samples, + &string_ids, + )?, + TypedSamples::Boolean(samples) => { + publish_boolean_values(appenders.get("boolean_values")?, sensor_id, samples)? + } + TypedSamples::Location(samples) => { + publish_location_values(appenders.get("location_values")?, sensor_id, samples)? + } + TypedSamples::Blob(samples) => { + publish_blob_values(appenders.get("blob_values")?, sensor_id, samples)? + } + TypedSamples::Json(samples) => { + publish_json_values(appenders.get("json_values")?, sensor_id, samples)? + } + } + } + appenders.flush() +} + +/// The appenders of a batch, opened when a sensor of that table shows up. +struct Appenders<'a> { + connection: &'a Connection, + appenders: HashMap<&'static str, Appender<'a>>, +} + +impl<'a> Appenders<'a> { + fn new(connection: &'a Connection) -> Self { + Self { + connection, + appenders: HashMap::new(), + } + } + + fn get(&mut self, table: &'static str) -> Result<&mut Appender<'a>> { + if !self.appenders.contains_key(table) { + self.appenders + .insert(table, self.connection.appender(table)?); + } + Ok(self.appenders.get_mut(table).expect("inserted above")) + } + + fn flush(mut self) -> Result<()> { + for appender in self.appenders.values_mut() { + appender.flush()?; + } + Ok(()) + } +} + +fn publish_integer_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample], ) -> Result<()> { - let mut appender = transaction.appender("integer_values")?; for value in values { let timestamp_us = value.datetime.to_rfc3339(); appender.append_row(params![sensor_id, timestamp_us, value.value])?; } - appender.flush()?; Ok(()) } -pub fn publish_numeric_values( - transaction: &Transaction, +fn publish_numeric_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample], ) -> Result<()> { - let mut appender = transaction.appender("numeric_values")?; for value in values { let timestamp_us = value.datetime.to_rfc3339(); let string_value = value.value.to_string(); appender.append_row(params![sensor_id, timestamp_us, string_value])?; } - appender.flush()?; Ok(()) } -pub fn publish_float_values( - transaction: &Transaction, +fn publish_float_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample], ) -> Result<()> { - let mut appender = transaction.appender("float_values")?; for value in values { let timestamp_us = value.datetime.to_rfc3339(); appender.append_row(params![sensor_id, timestamp_us, value.value])?; } - appender.flush()?; Ok(()) } -pub fn publish_string_values( - transaction: &Transaction, +fn publish_string_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample], + string_ids: &HashMap, ) -> Result<()> { - let mut appender = transaction.appender("string_values")?; for value in values { - let string_id = get_string_value_id_or_create(transaction, &value.value)?; + let string_id = string_ids + .get(&value.value) + .context("a registered string has no id")?; let timestamp_us = value.datetime.to_rfc3339(); appender.append_row(params![sensor_id, timestamp_us, string_id])?; } - appender.flush()?; Ok(()) } -pub fn publish_boolean_values( - transaction: &Transaction, +fn publish_boolean_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample], ) -> Result<()> { - let mut appender = transaction.appender("boolean_values")?; for value in values { let timestamp_us = value.datetime.to_rfc3339(); appender.append_row(params![sensor_id, timestamp_us, value.value])?; } - appender.flush()?; Ok(()) } -pub fn publish_location_values( - transaction: &Transaction, +fn publish_location_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample], ) -> Result<()> { - let mut appender = transaction.appender("location_values")?; for value in values { let timestamp_us = value.datetime.to_rfc3339(); let lat = value.value.y(); let lon = value.value.x(); appender.append_row(params![sensor_id, timestamp_us, lat, lon])?; } - appender.flush()?; Ok(()) } -pub fn publish_blob_values( - transaction: &Transaction, +fn publish_blob_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample>], ) -> Result<()> { - let mut appender = transaction.appender("blob_values")?; for value in values { let timestamp_us = value.datetime.to_rfc3339(); appender.append_row(params![sensor_id, timestamp_us, &value.value])?; } - appender.flush()?; Ok(()) } -pub fn publish_json_values( - transaction: &Transaction, +fn publish_json_values( + appender: &mut Appender, sensor_id: i64, values: &[Sample], ) -> Result<()> { - let mut appender = transaction.appender("json_values")?; for value in values { let timestamp_us = value.datetime.to_rfc3339(); let string_value = value.value.to_string(); appender.append_row(params![sensor_id, timestamp_us, string_value])?; } - appender.flush()?; Ok(()) } diff --git a/src/storage/duckdb/duckdb_registration.rs b/src/storage/duckdb/duckdb_registration.rs new file mode 100644 index 00000000..689763ce --- /dev/null +++ b/src/storage/duckdb/duckdb_registration.rs @@ -0,0 +1,388 @@ +//! Registering the sensors and the strings of a batch on DuckDB with a handful of statements, +//! whatever the number of sensors. +//! +//! A Prometheus request carries thousands of series. Registering them one by one (a lookup, an +//! insert, then a lookup and an insert per label) took about 0.9 ms per new series. The ids are +//! looked up with chunked `IN` lists, and the missing rows go in through appenders. Appenders take +//! the ids of the dictionaries and of the sensors from the sequences of the tables, so only the +//! columns that matter are given. +//! +//! Everything runs in the transaction of the batch, and nothing is cached: a failed batch leaves +//! nothing behind, and a deleted sensor cannot leave a stale id. The connection is behind a mutex, +//! so no other writer of this process runs at the same time. + +use crate::datamodel::Sensor; +use anyhow::{Result, bail}; +use duckdb::{Connection, params, params_from_iter}; +use std::collections::{BTreeSet, HashMap}; +use uuid::Uuid; + +/// How many keys one `IN (...)` lookup carries. +const LOOKUP_CHUNK: usize = 500; + +/// The sensor ids of the given sensors, creating the ones that do not exist. +/// +/// The labels of an existing sensor are left as they are: they are written when the sensor is +/// created. +pub fn register_sensors( + connection: &Connection, + sensors: &[&Sensor], +) -> Result> { + let mut unique: HashMap = HashMap::with_capacity(sensors.len()); + for sensor in sensors { + unique.entry(sensor.uuid).or_insert(sensor); + } + + let uuids: Vec = unique.keys().map(Uuid::to_string).collect(); + let mut ids: HashMap = + fetch_ids(connection, "sensors", "uuid", "sensor_id", &uuids)? + .into_iter() + .map(|(uuid, id)| Ok((Uuid::parse_str(&uuid)?, id))) + .collect::>()?; + + let mut missing: Vec<&Sensor> = unique + .values() + .filter(|sensor| !ids.contains_key(&sensor.uuid)) + .copied() + .collect(); + if missing.is_empty() { + return Ok(ids); + } + missing.sort_by_key(|sensor| sensor.uuid); + + let unit_ids = ensure_units(connection, &missing)?; + let label_names = ensure_dictionary( + connection, + "labels_name_dictionary", + "name", + &missing + .iter() + .flat_map(|sensor| sensor.labels.iter().map(|(name, _)| name.as_str())) + .collect(), + )?; + let label_descriptions = ensure_dictionary( + connection, + "labels_description_dictionary", + "description", + &missing + .iter() + .flat_map(|sensor| { + sensor + .labels + .iter() + .map(|(_, description)| description.as_str()) + }) + .collect(), + )?; + + { + let mut appender = connection.appender("sensors")?; + for column in ["uuid", "name", "type", "unit"] { + appender.add_column(column)?; + } + for sensor in &missing { + let unit_id: Option = sensor + .unit + .as_ref() + .and_then(|unit| unit_ids.get(&unit.name).copied()); + appender.append_row(params![ + sensor.uuid.to_string(), + sensor.name, + sensor.sensor_type.to_string(), + unit_id + ])?; + } + appender.flush()?; + } + + let new_uuids: Vec = missing + .iter() + .map(|sensor| sensor.uuid.to_string()) + .collect(); + for (uuid, id) in fetch_ids(connection, "sensors", "uuid", "sensor_id", &new_uuids)? { + ids.insert(Uuid::parse_str(&uuid)?, id); + } + + { + let mut appender = connection.appender("labels")?; + for sensor in &missing { + let Some(sensor_id) = ids.get(&sensor.uuid) else { + bail!("sensor {} was not created", sensor.uuid); + }; + for (name, description) in sensor.labels.iter() { + appender.append_row(params![ + sensor_id, + label_names[name.as_str()], + label_descriptions[description.as_str()] + ])?; + } + } + appender.flush()?; + } + + Ok(ids) +} + +/// The dictionary ids of the given strings, adding the ones that are not in it yet. +pub fn ensure_string_ids( + connection: &Connection, + strings: &BTreeSet<&str>, +) -> Result> { + ensure_dictionary(connection, "strings_values_dictionary", "value", strings) +} + +/// The ids of the units of the given sensors, adding the ones that do not exist. The description +/// of a unit that exists already is kept. +fn ensure_units(connection: &Connection, sensors: &[&Sensor]) -> Result> { + let mut units: HashMap<&str, Option<&str>> = HashMap::new(); + for sensor in sensors { + if let Some(unit) = &sensor.unit { + units + .entry(unit.name.as_str()) + .or_insert(unit.description.as_deref()); + } + } + if units.is_empty() { + return Ok(HashMap::new()); + } + + let names: Vec = units.keys().map(|name| name.to_string()).collect(); + let mut ids = fetch_ids(connection, "units", "name", "id", &names)?; + let mut missing: Vec<(&str, Option<&str>)> = units + .into_iter() + .filter(|(name, _)| !ids.contains_key(*name)) + .collect(); + if missing.is_empty() { + return Ok(ids); + } + missing.sort(); + + { + let mut appender = connection.appender("units")?; + appender.add_column("name")?; + appender.add_column("description")?; + for (name, description) in &missing { + appender.append_row(params![name, description])?; + } + appender.flush()?; + } + let new_names: Vec = missing.iter().map(|(name, _)| name.to_string()).collect(); + ids.extend(fetch_ids(connection, "units", "name", "id", &new_names)?); + Ok(ids) +} + +/// The ids of the given values of a dictionary table with one text column, adding the missing +/// ones. The table and the column are static names, never something a caller chooses. +fn ensure_dictionary( + connection: &Connection, + table: &'static str, + column: &'static str, + values: &BTreeSet<&str>, +) -> Result> { + if values.is_empty() { + return Ok(HashMap::new()); + } + let values: Vec = values.iter().map(|value| value.to_string()).collect(); + let mut ids = fetch_ids(connection, table, column, "id", &values)?; + let missing: Vec<&String> = values + .iter() + .filter(|value| !ids.contains_key(*value)) + .collect(); + if missing.is_empty() { + return Ok(ids); + } + + { + let mut appender = connection.appender(table)?; + appender.add_column(column)?; + for value in &missing { + appender.append_row(params![value])?; + } + appender.flush()?; + } + let new_values: Vec = missing.into_iter().cloned().collect(); + ids.extend(fetch_ids(connection, table, column, "id", &new_values)?); + Ok(ids) +} + +/// `key -> id` for the rows of `table` whose `key_column` is one of the keys. The table and the +/// columns are static names; the keys are bound as parameters, a chunk at a time. +fn fetch_ids( + connection: &Connection, + table: &'static str, + key_column: &'static str, + id_column: &'static str, + keys: &[String], +) -> Result> { + let mut ids = HashMap::with_capacity(keys.len()); + for chunk in keys.chunks(LOOKUP_CHUNK) { + let placeholders = vec!["?"; chunk.len()].join(", "); + let sql = format!( + "SELECT CAST({key_column} AS VARCHAR), {id_column} FROM {table} \ + WHERE {key_column} IN ({placeholders})" + ); + let mut statement = connection.prepare(&sql)?; + let rows = statement.query_map(params_from_iter(chunk.iter()), |row| { + Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?)) + })?; + for row in rows { + let (key, id) = row?; + ids.insert(key, id); + } + } + Ok(ids) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::datamodel::SensorType; + use crate::datamodel::sensapp_vec::SensAppLabels; + use crate::datamodel::unit::Unit; + + fn connection() -> Connection { + let connection = Connection::open_in_memory().unwrap(); + connection.execute_batch(super::super::INIT_SQL).unwrap(); + connection + } + + fn sensor(name: &str, unit: Option, labels: &[(&str, &str)]) -> Sensor { + let labels: SensAppLabels = labels + .iter() + .map(|(name, value)| (name.to_string(), value.to_string())) + .collect(); + Sensor::new( + Uuid::new_v4(), + name.to_string(), + SensorType::Float, + unit, + Some(labels), + ) + } + + fn count(connection: &Connection, table: &str) -> i64 { + connection + .query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| { + row.get(0) + }) + .unwrap() + } + + #[test] + fn registers_sensors_with_their_units_and_labels_once() { + let connection = connection(); + let celsius = Some(Unit::new("Cel".to_string(), Some("degrees".to_string()))); + let first = sensor("a", celsius.clone(), &[("room", "kitchen"), ("floor", "1")]); + let second = sensor("b", celsius, &[("room", "kitchen")]); + + // The same sensor twice in the input is one sensor + let ids = register_sensors(&connection, &[&first, &second, &first]).unwrap(); + assert_eq!(ids.len(), 2); + assert_ne!(ids[&first.uuid], ids[&second.uuid]); + assert_eq!(count(&connection, "sensors"), 2); + assert_eq!(count(&connection, "units"), 1); + assert_eq!(count(&connection, "labels"), 3); + assert_eq!(count(&connection, "labels_name_dictionary"), 2); + assert_eq!(count(&connection, "labels_description_dictionary"), 2); + + // Registering again finds the same ids and writes nothing + let again = register_sensors(&connection, &[&second, &first]).unwrap(); + assert_eq!(again, ids); + assert_eq!(count(&connection, "sensors"), 2); + assert_eq!(count(&connection, "labels"), 3); + + // A unit that exists keeps its description + let renamed = sensor( + "c", + Some(Unit::new("Cel".to_string(), Some("other".to_string()))), + &[], + ); + register_sensors(&connection, &[&renamed]).unwrap(); + let description: String = connection + .query_row( + "SELECT description FROM units WHERE name = 'Cel'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(description, "degrees"); + } + + #[test] + fn a_rolled_back_batch_leaves_nothing_and_can_be_written_again() { + let mut connection = connection(); + let kelvin = Some(Unit::new("K".to_string(), None)); + let sensors: Vec = (0..3) + .map(|index| sensor(&format!("s{index}"), kelvin.clone(), &[("room", "lab")])) + .collect(); + let refs: Vec<&Sensor> = sensors.iter().collect(); + + let transaction = connection.transaction().unwrap(); + register_sensors(&transaction, &refs).unwrap(); + ensure_string_ids(&transaction, &BTreeSet::from(["on", "off"])).unwrap(); + transaction.rollback().unwrap(); + for table in [ + "sensors", + "units", + "labels", + "labels_name_dictionary", + "labels_description_dictionary", + "strings_values_dictionary", + ] { + assert_eq!(count(&connection, table), 0, "{table}"); + } + + let transaction = connection.transaction().unwrap(); + let ids = register_sensors(&transaction, &refs).unwrap(); + transaction.commit().unwrap(); + assert_eq!(ids.len(), 3); + assert_eq!(count(&connection, "sensors"), 3); + assert_eq!(count(&connection, "labels"), 3); + assert_eq!(count(&connection, "units"), 1); + } + + #[test] + fn lookups_cross_the_chunk_size() { + let connection = connection(); + let total = LOOKUP_CHUNK * 2 + 17; + let sensors: Vec = (0..total) + .map(|index| sensor(&format!("s{index}"), None, &[("n", &format!("v{index}"))])) + .collect(); + let refs: Vec<&Sensor> = sensors.iter().collect(); + + let ids = register_sensors(&connection, &refs).unwrap(); + assert_eq!(ids.len(), total); + assert_eq!(count(&connection, "labels"), total as i64); + assert_eq!(register_sensors(&connection, &refs).unwrap(), ids); + assert_eq!(count(&connection, "sensors"), total as i64); + + let owned: Vec = (0..total) + .map(|index| format!("string {index} é 日本")) + .collect(); + let strings: BTreeSet<&str> = owned.iter().map(String::as_str).collect(); + let string_ids = ensure_string_ids(&connection, &strings).unwrap(); + assert_eq!(string_ids.len(), total); + assert_eq!( + ensure_string_ids(&connection, &strings).unwrap(), + string_ids + ); + assert_eq!( + count(&connection, "strings_values_dictionary"), + total as i64 + ); + } + + #[test] + fn empty_input_registers_nothing() { + let connection = connection(); + assert!(register_sensors(&connection, &[]).unwrap().is_empty()); + assert!( + ensure_string_ids(&connection, &BTreeSet::new()) + .unwrap() + .is_empty() + ); + let unlabelled = sensor("a", None, &[]); + register_sensors(&connection, &[&unlabelled]).unwrap(); + assert_eq!(count(&connection, "labels"), 0); + } +} diff --git a/src/storage/duckdb/duckdb_utilities.rs b/src/storage/duckdb/duckdb_utilities.rs deleted file mode 100644 index d874d17f..00000000 --- a/src/storage/duckdb/duckdb_utilities.rs +++ /dev/null @@ -1,167 +0,0 @@ -use crate::datamodel::Sensor; -use crate::datamodel::unit::Unit; -use anyhow::Result; -use cached::macros::cached; -use duckdb::{CachedStatement, OptionalExt, Transaction, params}; -use uuid::Uuid; - -#[cached( - ttl_secs = 120, - sync_writes = "default", - key = "String", - convert = { label_name.to_string() } -)] -pub fn get_label_name_id_or_create(transaction: &Transaction, label_name: &str) -> Result { - let mut select_stmt: CachedStatement = - transaction.prepare_cached("SELECT id FROM labels_name_dictionary WHERE name = ?")?; - - let label_name_id = select_stmt - .query_row(params![label_name], |row| row.get(0)) - .optional()?; - - if let Some(id) = label_name_id { - Ok(id) - } else { - let mut insert_stmt: CachedStatement = transaction - .prepare_cached("INSERT INTO labels_name_dictionary (name) VALUES (?) RETURNING id")?; - let label_name_id: i64 = insert_stmt.query_row(params![label_name], |row| row.get(0))?; - Ok(label_name_id) - } -} - -#[cached( - ttl_secs = 120, - sync_writes = "default", - key = "String", - convert = { label_description.to_string() } -)] -pub fn get_label_description_id_or_create( - transaction: &Transaction, - label_description: &str, -) -> Result { - let mut select_stmt: CachedStatement = transaction - .prepare_cached("SELECT id FROM labels_description_dictionary WHERE description = ?")?; - - let label_description_id = select_stmt - .query_row(params![label_description], |row| row.get(0)) - .optional()?; - - if let Some(id) = label_description_id { - Ok(id) - } else { - let mut insert_stmt: CachedStatement = transaction.prepare_cached( - "INSERT INTO labels_description_dictionary (description) VALUES (?) RETURNING id", - )?; - let label_description_id: i64 = - insert_stmt.query_row(params![label_description], |row| row.get(0))?; - Ok(label_description_id) - } -} - -#[cached( - ttl_secs = 120, - sync_writes = "default", - key = "String", - convert = { unit.name.clone() } -)] -pub fn get_unit_id_or_create(transaction: &Transaction, unit: &Unit) -> Result { - let mut select_stmt: CachedStatement = - transaction.prepare_cached("SELECT id FROM units WHERE name = ?")?; - - let unit_id = select_stmt - .query_row(params![unit.name], |row| row.get(0)) - .optional()?; - - if let Some(id) = unit_id { - Ok(id) - } else { - let mut insert_stmt: CachedStatement = transaction - .prepare_cached("INSERT INTO units (name, description) VALUES (?, ?) RETURNING id")?; - let unit_id: i64 = - insert_stmt.query_row(params![unit.name, unit.description], |row| row.get(0))?; - Ok(unit_id) - } -} - -#[cached( - ttl_secs = 120, - sync_writes = "default", - key = "Uuid", - convert = { sensor.uuid } -)] -pub fn get_sensor_id_or_create_sensor(transaction: &Transaction, sensor: &Sensor) -> Result { - let uuid_string = sensor.uuid.to_string(); - - let mut select_stmt: CachedStatement = - transaction.prepare_cached("SELECT sensor_id FROM sensors WHERE uuid = ?")?; - - let existing_sensor_id = select_stmt - .query_row(params![uuid_string], |row| row.get(0)) - .optional()?; - - if let Some(existing_sensor_id) = existing_sensor_id { - Ok(existing_sensor_id) - } else { - let sensor_type_string = sensor.sensor_type.to_string(); - - let unit_id = match sensor.unit { - Some(ref unit) => Some(get_unit_id_or_create(transaction, unit)?), - None => None, - }; - - let mut insert_stmt: CachedStatement = transaction.prepare_cached( - "INSERT INTO sensors (uuid, name, type, unit) VALUES (?, ?, ?, ?) RETURNING sensor_id", - )?; - let sensor_id: i64 = insert_stmt.query_row( - params![uuid_string, sensor.name, sensor_type_string, unit_id], - |row| row.get(0), - )?; - - // Add the labels - let mut label_insert_stmt: CachedStatement = transaction - .prepare_cached("INSERT INTO labels (sensor_id, name, description) VALUES (?, ?, ?)")?; - for (key, value) in sensor.labels.iter() { - let label_name_id = get_label_name_id_or_create(transaction, key)?; - let label_description_id = get_label_description_id_or_create(transaction, value)?; - label_insert_stmt.execute(params![sensor_id, label_name_id, label_description_id])?; - } - - Ok(sensor_id) - } -} - -#[cached( - ttl_secs = 120, - sync_writes = "default", - key = "String", - convert = { string_value.to_string() } -)] -pub fn get_string_value_id_or_create(transaction: &Transaction, string_value: &str) -> Result { - let mut select_stmt: CachedStatement = - transaction.prepare_cached("SELECT id FROM strings_values_dictionary WHERE value = ?")?; - - let string_id = select_stmt - .query_row(params![string_value], |row| row.get(0)) - .optional()?; - - if let Some(id) = string_id { - Ok(id) - } else { - let mut insert_stmt: CachedStatement = transaction.prepare_cached( - "INSERT INTO strings_values_dictionary (value) VALUES (?) RETURNING id", - )?; - let string_id: i64 = insert_stmt.query_row(params![string_value], |row| row.get(0))?; - Ok(string_id) - } -} - -/// Forget the cached id of a deleted sensor. -/// -/// Without this, publishing the same sensor again within the cache TTL would write -/// samples under a sensor id that no longer exists. -pub fn forget_sensor_id(sensor_uuid: &Uuid) { - use cached::Cached; - let _ = GET_SENSOR_ID_OR_CREATE_SENSOR - .write() - .cache_remove(sensor_uuid); -} diff --git a/src/storage/duckdb/mod.rs b/src/storage/duckdb/mod.rs index 2553c6f6..5d5e3bb2 100644 --- a/src/storage/duckdb/mod.rs +++ b/src/storage/duckdb/mod.rs @@ -1,4 +1,4 @@ -use crate::datamodel::batch::{Batch, SingleSensorBatch}; +use crate::datamodel::batch::Batch; use crate::datamodel::sensapp_datetime::SensAppDateTimeExt; use crate::datamodel::sensapp_vec::SensAppLabels; use crate::datamodel::unit::Unit; @@ -8,8 +8,7 @@ use crate::datamodel::{ use anyhow::{Context, Result, bail}; use async_trait::async_trait; use duckdb::{Connection, OptionalExt}; -use duckdb_publishers::*; -use duckdb_utilities::{forget_sensor_id, get_sensor_id_or_create_sensor}; +use duckdb_publishers::publish_batch; use geo::Point; use rust_decimal::Decimal; use serde_json::Value as JsonValue; @@ -26,7 +25,7 @@ use super::{ }; mod duckdb_publishers; -mod duckdb_utilities; +mod duckdb_registration; mod selector; #[derive(Debug)] @@ -66,9 +65,7 @@ impl StorageInstance for DuckDBStorage { spawn_blocking(move || -> Result<()> { let mut connection = connection.blocking_lock(); let transaction = connection.transaction()?; - for single_sensor_batch in bbatch.sensors.as_ref() { - publish_single_sensor_batch(&transaction, single_sensor_batch)?; - } + publish_batch(&transaction, bbatch.sensors.as_ref())?; transaction.commit()?; Ok(()) }) @@ -134,7 +131,6 @@ impl StorageInstance for DuckDBStorage { transaction.commit()?; connection.execute("DELETE FROM sensors WHERE sensor_id = ?", [sensor_id])?; - forget_sensor_id(&parsed_uuid); Ok(true) }) .await? @@ -979,66 +975,10 @@ impl StorageInstance for DuckDBStorage { .ok(); connection.execute("DELETE FROM units", []).ok(); - // Step 2: Clear all cached function caches - // The cached macro generates cache variables named after the function in uppercase - use cached::Cached; - duckdb_utilities::GET_LABEL_NAME_ID_OR_CREATE - .write() - .cache_clear(); - duckdb_utilities::GET_LABEL_DESCRIPTION_ID_OR_CREATE - .write() - .cache_clear(); - duckdb_utilities::GET_UNIT_ID_OR_CREATE - .write() - .cache_clear(); - duckdb_utilities::GET_SENSOR_ID_OR_CREATE_SENSOR - .write() - .cache_clear(); - duckdb_utilities::GET_STRING_VALUE_ID_OR_CREATE - .write() - .cache_clear(); - Ok(()) } } -fn publish_single_sensor_batch( - transaction: &duckdb::Transaction, - single_sensor_batch: &SingleSensorBatch, -) -> Result<()> { - let sensor_id = get_sensor_id_or_create_sensor(transaction, &single_sensor_batch.sensor)?; - { - let samples_guard = single_sensor_batch.samples.blocking_read(); - match &*samples_guard { - TypedSamples::Integer(samples) => { - publish_integer_values(transaction, sensor_id, samples)?; - } - TypedSamples::Numeric(samples) => { - publish_numeric_values(transaction, sensor_id, samples)?; - } - TypedSamples::Float(samples) => { - publish_float_values(transaction, sensor_id, samples)?; - } - TypedSamples::String(samples) => { - publish_string_values(transaction, sensor_id, samples)?; - } - TypedSamples::Boolean(samples) => { - publish_boolean_values(transaction, sensor_id, samples)?; - } - TypedSamples::Location(samples) => { - publish_location_values(transaction, sensor_id, samples)?; - } - TypedSamples::Blob(samples) => { - publish_blob_values(transaction, sensor_id, samples)?; - } - TypedSamples::Json(samples) => { - publish_json_values(transaction, sensor_id, samples)?; - } - } - } - Ok(()) -} - fn duckdb_get_sensor_metadata( connection: &Connection, sensor_uuid: &str, From c06774b7f24f713094024b523798264e7e79338e Mon Sep 17 00:00:00 2001 From: Antoine Pultier Date: Fri, 2 Oct 2026 22:53:17 +0200 Subject: [PATCH 4/4] =?UTF-8?q?docs:=20=F0=9F=93=9D=20close=20the=20DuckDB?= =?UTF-8?q?=20bulk=20writes=20task,=20with=20the=20numbers?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Also correct the backends page: DuckDB stores microsecond timestamps since the TIMESTAMP_MS columns were replaced. Co-Authored-By: Claude Sonnet 5.5 --- current_tasks/duckdb-bulk-writes.md | 46 ---------------- docs/BACKENDS.md | 4 +- done/duckdb-bulk-writes.md | 82 +++++++++++++++++++++++++++++ 3 files changed, 84 insertions(+), 48 deletions(-) delete mode 100644 current_tasks/duckdb-bulk-writes.md create mode 100644 done/duckdb-bulk-writes.md diff --git a/current_tasks/duckdb-bulk-writes.md b/current_tasks/duckdb-bulk-writes.md deleted file mode 100644 index 480a9e1b..00000000 --- a/current_tasks/duckdb-bulk-writes.md +++ /dev/null @@ -1,46 +0,0 @@ -# DuckDB: bulk writes - -## Goal - -Writing many series, or many strings, to DuckDB must cost a handful of statements whatever the number of -series, like PostgreSQL, TimescaleDB and ClickHouse. DuckDB is not the main target (local analysis), but the -per-series path is slow and the `#[cached]` ids have the rollback problem the other backends lost. - -## Baseline (2 Oct 2026, release build, macOS, `tests/perf/scale.sh 3000` and `tests/perf/strings.sh 3000`) - -| write | before | -|---|---| -| 3000 new series (30 000 samples) | 2.69 s | -| the same series again | 0.37 s | -| strings: 3000 series x 10, 50 distinct strings | 2.51 s | -| strings: one series of 10 000 distinct strings | 1.50 s | -| strings: the first request again | 0.37 s | - -Why: `get_sensor_id_or_create_sensor` is one `SELECT`, one `INSERT` and one `INSERT` per label (plus the -dictionaries) per sensor, the string dictionary is one statement per distinct string, and every sensor opens -its own appender on its value table. All of it behind `#[cached]` wrappers that are filled inside the -transaction, so a rollback leaves ids of sensors that were never committed (until the 120 s TTL ends). - -## Plan - -1. `duckdb_registration.rs`: register all the sensors of a batch with a few statements (ids by chunked - `IN`, appenders for units, label dictionaries, sensors and labels). Nothing cached. -2. Bulk string dictionary for the whole batch. -3. One appender per value table for the whole batch instead of one per sensor. -4. Remove the `#[cached]` wrappers, `forget_sensor_id` and their cache clears. -5. Before/after numbers, tests on top of the backend-generic ones, docs. - -## Done when - -- [ ] Before/after numbers below. -- [ ] A failed batch leaves no sensor, label or string behind and a later write of the same series works. -- [ ] Backend-generic tests still pass on DuckDB, and on SQLite and TimescaleDB where they are generic. -- [ ] Full suites, clippy, and the DuckDB suite. - -## Not in scope - -The aggregated selector read (`BulkSelectorBackend::read_aggregated_samples`) is still one query per series on -DuckDB (0.30 s for a 100 series remote read with a step). The `time_bucket` SQL of `duckdb_bucketed_cte` -extends to many sensors with `GROUP BY sensor_id, bucket`. A read, not a write: separate task if wanted. - -## Progress diff --git a/docs/BACKENDS.md b/docs/BACKENDS.md index dc15f0ba..e7a3f833 100644 --- a/docs/BACKENDS.md +++ b/docs/BACKENDS.md @@ -12,7 +12,7 @@ mature. This page says which ones to rely on, and what each one is for. | PostgreSQL | **Maintained**, main development backend | yes | Small and medium deployments, development | | TimescaleDB | **Maintained** | yes | PostgreSQL deployments that want hypertables and compression | | SQLite | **Maintained** | yes | Tests, demos, single-node edge devices, one writer | -| DuckDB | Compatibility path, **less mature** | yes | Local analysis of a dataset, notebooks. Stores millisecond timestamps | +| DuckDB | Compatibility path, **less mature** | yes | Local analysis of a dataset, notebooks. Writes are bulk (3000 new series in 0.2 s) | | RRDCached | **Experimental** | yes (its own job) | Fixed-size round-robin storage, monitoring-style data; no deletion | | BigQuery | **Experimental, parked** | compile only | Nothing yet: it has not been brought back in line with the current storage interface (see `ideas/bigquery-backend-reconciliation.md`) | @@ -52,7 +52,7 @@ or writes series one by one still works, it costs more round trips with many ser returned wrong rows after a chunk changed state, and the aggregated `count` is written `COUNT(value)`, because `COUNT(*)` failed to plan over many chunks. Details: `done/timescaledb-compressed-chunks.md`. - **SQLite**: one writer at a time; units keep their name but not their description. -- **DuckDB**: timestamps are stored with a millisecond precision, unlike the other backends (microseconds). +- **DuckDB**: timestamps are stored with a microsecond precision, like the other backends. One process owns the database file; a batch is one transaction, registered and written with a few statements whatever the number of series. - **RRDCached**: only its dedicated integration module runs against a real `rrdcached` (the backend-generic suite does not apply: no labels, no deletion, consolidated data). Data is consolidated according to the chosen preset, old precision is lost by design. Details: [RRDCACHED.md](RRDCACHED.md). diff --git a/done/duckdb-bulk-writes.md b/done/duckdb-bulk-writes.md new file mode 100644 index 00000000..549badba --- /dev/null +++ b/done/duckdb-bulk-writes.md @@ -0,0 +1,82 @@ +# DuckDB: bulk writes + +## Goal + +Writing many series, or many strings, to DuckDB must cost a handful of statements whatever the number of +series, like PostgreSQL, TimescaleDB and ClickHouse. DuckDB is not the main target (local analysis), but the +per-series path is slow and the `#[cached]` ids have the rollback problem the other backends lost. + +## Baseline (2 Oct 2026, release build, macOS, `tests/perf/scale.sh 3000` and `tests/perf/strings.sh 3000`) + +| write | before | +|---|---| +| 3000 new series (30 000 samples) | 2.69 s | +| the same series again | 0.37 s | +| strings: 3000 series x 10, 50 distinct strings | 2.51 s | +| strings: one series of 10 000 distinct strings | 1.50 s | +| strings: the first request again | 0.37 s | + +Why: `get_sensor_id_or_create_sensor` is one `SELECT`, one `INSERT` and one `INSERT` per label (plus the +dictionaries) per sensor, the string dictionary is one statement per distinct string, and every sensor opens +its own appender on its value table. All of it behind `#[cached]` wrappers that are filled inside the +transaction, so a rollback leaves ids of sensors that were never committed (until the 120 s TTL ends). + +## Plan + +1. `duckdb_registration.rs`: register all the sensors of a batch with a few statements (ids by chunked + `IN`, appenders for units, label dictionaries, sensors and labels). Nothing cached. +2. Bulk string dictionary for the whole batch. +3. One appender per value table for the whole batch instead of one per sensor. +4. Remove the `#[cached]` wrappers, `forget_sensor_id` and their cache clears. +5. Before/after numbers, tests on top of the backend-generic ones, docs. + +## Done when + +- [x] Before/after numbers below. +- [x] A failed batch leaves no sensor, label or string behind and a later write of the same series works. +- [x] Backend-generic tests still pass on DuckDB, and on SQLite and TimescaleDB where they are generic. +- [x] Full suites, clippy, and the DuckDB suite. + +## Not in scope + +The aggregated selector read (`BulkSelectorBackend::read_aggregated_samples`) is still one query per series on +DuckDB (0.30 s for a 100 series remote read with a step). The `time_bucket` SQL of `duckdb_bucketed_cte` +extends to many sensors with `GROUP BY sensor_id, bucket`. A read, not a write: separate task if wanted. + +## Progress + +Done 2 Oct 2026. + +- New `src/storage/duckdb/duckdb_registration.rs`: `register_sensors` and `ensure_string_ids`. The ids that exist + are read with chunked `IN` lists (the crate cannot bind list parameters), the missing units, label names and + descriptions, strings, sensors and labels go in through appenders with `add_column`, so the ids still come + from the sequences of the tables. The same unit, label or string twice in a batch is one row; a unit that + exists keeps its description; labels are written when the sensor is created, like on the other backends. +- `duckdb_publishers.rs`: `publish_batch` registers everything first, then writes the samples of every sensor + through one appender per value table. The string samples use the ids of the batch dictionary. +- `duckdb_utilities.rs` is gone with its `#[cached]` wrappers, `forget_sensor_id` and the cache clears of + `delete_series` and of the test cleanup. A rolled back batch cannot leave a stale sensor id any more. +- Tests: four unit tests on an in-memory database (one unit, labels and dictionaries written once, the same sensor + twice in the input, a rollback leaves every table empty and the same batch can be written again, lookups and + strings across the chunk size, empty input). The backend-generic tests already cover 2100 new sensors in one + batch, string dictionaries, labels written once and concurrent first writes; the DuckDB suite (236 lib, 274 + integration) and the default build (239 lib, 287 integration) pass, `clippy --tests -D warnings` on both. + +Numbers, same machine and scripts as the baseline (release build, DuckDB file, one run each, indicative): + +| write | before | after | +|---|---|---| +| 3000 new series (30 000 samples) | 2.69 s | 0.23 s | +| the same series again | 0.37 s | 0.11 s | +| strings: 3000 series x 10, 50 distinct strings | 2.51 s | 0.20 s | +| strings: one series of 10 000 distinct strings | 1.50 s | 0.16 s | +| strings: the first request again | 0.37 s | 0.10 s | + +Reads are unchanged (a 100 series remote read with a step is still about 0.3 s: one query per series). + +Not done: the timestamps still go through the appender as RFC 3339 text, parsed by DuckDB. At 3.6 us per sample +it is not what is left. + +Note for later: the macOS binary needs `libduckdb.dylib` next to it (`install_name_tool -add_rpath +@executable_path` and `codesign -f -s -` on a copy) to run the `tests/perf` scripts, and `cargo test --doc` cannot +find it locally. CI and the Docker image are not affected.