Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions docs/BACKENDS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`) |

Expand Down Expand Up @@ -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).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
82 changes: 82 additions & 0 deletions done/duckdb-bulk-writes.md
Original file line number Diff line number Diff line change
@@ -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.
16 changes: 0 additions & 16 deletions ideas/duckdb-bulk-writes.md

This file was deleted.

150 changes: 113 additions & 37 deletions src/storage/duckdb/duckdb_publishers.rs
Original file line number Diff line number Diff line change
@@ -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<i64>],
) -> 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<Decimal>],
) -> 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<f64>],
) -> 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>],
string_ids: &HashMap<String, i64>,
) -> 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<bool>],
) -> 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<Point>],
) -> 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<Vec<u8>>],
) -> 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<Value>],
) -> 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(())
}
Loading
Loading