# MeterStore
**Hot/cold tiered storage for metering time series.** PostgreSQL holds the recent
interval window at low latency; Apache Iceberg holds the history at analytical
scale. A single explicit timestamp separates them, and one SQL statement spans
both.
[](https://github.com/hupe1980/meterstore/actions/workflows/ci.yml)
[](#license)
π **[Documentation](https://hupe1980.github.io/meterstore)** Β· [API reference](https://docs.rs/meterstore) Β· [Changelog](CHANGELOG.md)
> **Pre-alpha, and unpublished on purpose.** Storage, tiering, archival,
> querying, reproducible reads and completeness work end to end against real
> PostgreSQL 16 and a real Iceberg warehouse. The API is still settling;
> integrating against a real workload is what settles it.
---
## The problem
An intelligent measuring system produces one value per measuring point, per OBIS
code, per interval. Fifteen minutes is the German settlement grain:
| 10 k measuring points | ~1 M | ~350 M |
| 100 k (mid-size utility) | ~9.6 M | ~3.5 B |
| 1 M (metering operator) | ~96 M | ~35 B |
Retention is regulatory β years to decades for the settlement record. PostgreSQL
handles the first row comfortably, the second with care, and the third not at all
without becoming a full-time job.
But the operational workload genuinely needs Postgres: recent data is written
continuously, corrected, and read transactionally by billing and market
communication. Meanwhile settlement, forecasting and grid analysis scan years
across hundreds of thousands of meters β an object-storage-and-columnar-format
problem.
**The data has a natural split most systems refuse to exploit:** recent intervals
are hot and still being corrected; historical intervals are cold and settled. The
boundary between them is a timestamp.
## How it works
```
MeterInterval.from ββββββββββββββββββββββββββββββββββββββΆ
ββββββ Iceberg (cold, settled) βββββΆβ
ββββ Postgres (hot) βββΆβ
epoch tiering_watermark now
```
A row's interval start alone decides its tier, so the tiers are disjoint by
construction β no deduplication, no merge, no double-counting. Four decisions
carry most of the weight:
- **The watermark lives inside the Iceberg snapshot.** Archival writes the tier
boundary into the snapshot summary in the same commit as the data. Iceberg
commits are a compare-and-swap, so rows and boundary become durable together or
not at all.
- **Purge is `DROP TABLE`, never `DELETE`.** The hot table is time-partitioned, so
archiving a window drops exactly one partition. Deleting a day of readings for
100 k meters row by row would leave ~9.6 M dead tuples for autovacuum.
- **Corrections are versions, not overwrites.** MSCONS corrects a value by
*versioning* it, so the store needs only Iceberg's `append` β and a past
settlement stays reproducible.
- **Nothing on the archival path holds a window.** Peak memory is the chunk size,
not the ~9.6 M-row window.
[How tiering works β](https://hupe1980.github.io/meterstore/docs/architecture/)
## Quick start
Without writing a program:
```bash
cargo install meterstore --features cli
meterstore init # a commented starter configuration
meterstore check # full validation β no database needed
meterstore create # both tiers, every declared table
meterstore status # boundary, lag, write runway, health
```
`meterstore query` runs SQL across both tiers and prints the boundary the answer
was computed against; `meterstore maintain` is the archival loop as a foreground
process. [The CLI β](https://hupe1980.github.io/meterstore/docs/cli/)
As a library:
```bash
cargo add meterstore
```
```rust
use meterstore::prelude::*;
use meterstore::hot::PostgresHot;
use std::sync::Arc;
use time::Duration;
let hot = Arc::new(PostgresHot::new(pool)); // a pool you already own
let cold = IcebergSqlCatalog {
database_url: &db_url,
warehouse_uri: "s3://bucket/warehouse", // or file:// memory:// gs:// abfss://
catalog_name: "meterstore",
namespace: "metering",
file_target_bytes: 512 * 1024 * 1024,
metadata_pool_max_connections: 4,
auth: &WarehouseAuth { region: Some("eu-central-1".into()), ..Default::default() },
}.build().await?.cold();
let store = MeterStore::builder()
.hot(hot)
.cold(cold.clone(), cold.table_provider("readings_versions").await?)
.table(
TableConfig::new("readings_versions")
.settlement_lag(Duration::days(7))
.archival_step(Duration::DAY)
.build()?,
)
.build()
.await?;
store.create_tables().await?;
```
Then write and read:
```rust
// Routes each interval to the tier that owns it.
store.append(&[stored_series]).await?;
// One statement, both tiers β with the boundary it was computed against.
let result = store.query(r#"
SELECT meter_local_day("from") AS day, SUM(value) AS kwh
FROM readings
WHERE malo_id = '41373559241'
AND "from" >= '2025-01-01' AND "from" < '2026-01-01'
GROUP BY 1 ORDER BY 1
"#).await?;
result.watermark(); // where cold ended and hot began
result.touched_hot_tier(); // whether the answer is only valid for now
```
A table declares whether it holds **spans or instants**. A *Lastgang* is energy
over `[from, to)`; a *ZΓ€hlerstandsgang* is a cumulative register value at an
instant, which BK6-24-174 (in force 06.06.2025) has made a primary record exactly
as voluminous as the Lastgang differenced out of it β and Β§ 146 Abs. 4 AO means
it cannot be discarded afterwards. Both tier the same way, because everything
that does the tiering reads the start timestamp:
```rust
TableConfig::new("meter_reads_versions").time_model(TimeModel::Point)
store.append_readings(&[zaehlerstandsgang]).await?; // metering::MeterReading
store.readings(malo)?.melo(melo)?.latest().await?; // what the meter reads now
```
They are never the same table: `value` is interval energy on one and a register
reading on the other, and summing the two together gives a number with no meaning
that looks exactly like a consumption total.
The shape also decides what **names** a reading. A Marktlokation may be measured
by several Messlokationen, and both meters carry `1-0:1.8.0` at the same instants.
A load profile belongs to the market location, so `melo_id` labels it; a register
belongs to the *meter*, so it joins the merge key β by default on a point table,
`identify_by_melo` either way. Keyed wrongly, two meters that agree on a number
store as one reading.
A `MeasurementSeries` holds one `obis_code`, so a typed read describes **one
channel** β folding import and export together sums to twice the truth. A
measuring point is a set of them, so there is a read for that too, in one scan:
```rust
store.series(malo)?.obis("1-0:1.8.0")?.range(from, to).collect().await?; // one channel
store.series(malo)?.range(from, to).collect_by_channel().await?; // all of them
```
`readings` is version-resolved; `readings_versions` is the raw audit trail. The
naming is load-bearing β see
[the version-resolution trap](https://hupe1980.github.io/meterstore/docs/interop/#the-version-resolution-trap)
before pointing an external engine at the warehouse.
A `SUM` over an incomplete month returns a smaller number and no reason, so
**completeness is a query**. The expected count is the DST-aware calendar's β 92
on the spring day, 100 on the autumn one, and for gas both on the Gastag rather
than the Sunday. The strongest finding is the one a range cannot make about
itself: a channel that delivered *nothing* produces no rows to aggregate, so it
needs a roster from an earlier window.
```rust
store.completeness(from, to).await?; // gaps within
store.completeness(from, to).seen_since(from - month).await?; // and what went silent
```
Personal data comes with a clock rather than a request. Β§ 60 Abs. 6 MsbG says
erase or anonymise *at the latest* three years after the end of the year a value
was collected in, so the sweep is a scheduled job β and a **catalogue** one,
because one subject map spans every table and a per-table sweep would orphan
readings the other tables still hold:
```rust
catalog.maintenance()
.anonymise_after(Retention::CalendarYears(3), "Β§ 60 Abs. 6 MsbG", "retention-job")
.spawn();
```
`meter_local_day` is not a convenience, and for **gas it is the wrong function**.
`Europe/Berlin` observes daylight saving, so the UTC day boundary sits at 01:00
or 02:00 local and grouping on UTC days is wrong every day of the year β but the
German gas market does not balance on the calendar day either. A *Gastag* runs
06:00 to 06:00 local, so a gas Lastgang grouped by the calendar day books six
hours a day into the neighbouring Bilanzierungstag, with totals that still look
plausible. `meter_balancing_day("from", sparte)` reads the commodity per row and
picks the right one:
```sql
SELECT sparte, meter_balancing_day("from", sparte) AS day, SUM(value)
FROM readings GROUP BY 1, 2;
```
The DST anomaly moves with the boundary: the clocks change *before* 06:00, so the
25-hour gas day is the one named after the **Saturday** while the 25-hour
calendar day is the Sunday.
And the boundary carries up to the **month**. The gas Bilanzierungsmonat runs
01.06 06:00 to 01.07 06:00 (EDI@Energy *Allgemeine Festlegungen* v6.1c, Kap. 3.1),
so an MSCONS version scope for a gas row is cut at 06:00 as well β which is why
every `VersionScope` constructor takes a `Sparte`:
```rust
// The network operator's Marktpartner-ID, parsed β because a wrong-but-plausible
// one is not an error, it is a different scope that nothing else shares.
VersionScope::for_interval("9900000000001", interval.from, Sparte::Gas)?
```
**An external engine does not get that function β so it gets the answer instead.**
SQL dialects differ on timestamp arithmetic, so no single published expression is
right everywhere. The encoder applies the calendar once, at write time, and stores
the answer:
```sql
-- Every engine. No zone conversion, no DST reasoning, no dialect.
SELECT balancing_day, SUM(value) FROM readings GROUP BY 1;
```
[Getting started β](https://hupe1980.github.io/meterstore/docs/getting-started/)
## Requirements
| Rust | 1.94 | Set by the dependency floor (`metering`, `iceberg`) |
| PostgreSQL | **12 or later** | `ATTACH PARTITION` takes only `SHARE UPDATE EXCLUSIVE` on the parent from 12 β see below |
| `metering` | **0.20 or later** | The domain layer β MeterStore stores its types, it does not redefine them |
| Apache Iceberg | format v2 | [Deliberately not v3](https://hupe1980.github.io/meterstore/docs/architecture/#format-version) |
Partition creation runs on the write path, and `CREATE TABLE β¦ PARTITION OF`
takes `ACCESS EXCLUSIVE` on the parent β which, since PostgreSQL grants locks in
arrival order, lets one long query stall every subsequent insert. Partitions are
built standalone and *attached* instead.
[Locks β](https://hupe1980.github.io/meterstore/docs/operations/#locks-and-why-ddl-gives-up)
The cold tier takes **any** `Arc<dyn Catalog>` β SQL, REST, Polaris, Lakekeeper,
Glue β and that seam is driven end to end by the test suite rather than asserted.
Three are built for you: a PostgreSQL-backed SQL catalogue on the same database as
the hot tier, a REST catalogue (`rest-catalog`, on by default), and AWS S3 Tables
behind the `s3tables` feature. A configuration file builds whichever it names β
`Settings::connect()` returns the pool, both tiers and every validated table.
[Details](https://hupe1980.github.io/meterstore/docs/getting-started/).
MeterStore needs only `SELECT` plus ownership of its own tables: no server
configuration, no restart, no extension. That is what makes it deployable on RDS,
Cloud SQL and Azure Postgres, where an extension-based approach is not.
## Relationship to `metering`
```
metering β what a measurement is, and how to compute with it (zero I/O, no async)
meterstore β where it lives, how it is tiered, how it is queried (all I/O)
```
[`metering`](https://crates.io/crates/metering) owns intervals, units, quality
flags, DST-correct calendars, the identifiers (`MaloId`, `MeloId`, `BdewCode`),
validation, Ersatzwertbildung, gas conversion and aggregation. MeterStore adds
exactly three things: **correction versioning**, the **transaction-time axis**,
and the **tiering boundary**.
That boundary is deliberate. Duplicating a domain rule here β a unit conversion, a
DST calendar β would create a second implementation to keep correct, and it would
drift.
## Documentation
| [Getting started](https://hupe1980.github.io/meterstore/docs/getting-started/) | Requirements, install, a store over both tiers |
| [Architecture](https://hupe1980.github.io/meterstore/docs/architecture/) | The watermark, the invariant, crash-safe archival |
| [Storage model](https://hupe1980.github.io/meterstore/docs/storage-model/) | Columns, the merge key, identity vs attribute, constraints |
| [Writing readings](https://hupe1980.github.io/meterstore/docs/writing/) | Routed writes, bulk ingest, idempotent redelivery |
| [Querying](https://hupe1980.github.io/meterstore/docs/querying/) | SQL across tiers, provenance, the typed series API |
| [Reproducibility](https://hupe1980.github.io/meterstore/docs/reproducibility/) | Settlement reruns on two independent time axes |
| [Completeness](https://hupe1980.github.io/meterstore/docs/completeness/) | DST-aware gap detection, including the channel that delivered nothing |
| [Operations](https://hupe1980.github.io/meterstore/docs/operations/) | Scheduling, locks, system tables, metrics, failure matrix |
| [The CLI](https://hupe1980.github.io/meterstore/docs/cli/) | `meterstore` β check, create, status, archive, maintain, query, serve |
| [External engines](https://hupe1980.github.io/meterstore/docs/interop/) | Spark, Trino, DuckDB β and the trap to avoid |
| [Privacy and retention](https://hupe1980.github.io/meterstore/docs/privacy/) | Pseudonymisation, and the three-year duty as a scheduled job |
| [Configuration](https://hupe1980.github.io/meterstore/docs/configuration/) | TOML over the same validated types |
## Status
Everything the documentation describes works end to end against real
infrastructure β both tiers, streaming archival, tier-split queries, reproducible
reads, completeness, multi-table sessions and both serving surfaces.
**807 tests**: unit, property, doc and integration against real PostgreSQL 16 and
a real Iceberg warehouse, plus an independently implemented correctness oracle over
generated workloads, covering both record shapes. **DuckDB** and **PyIceberg**
read the output and agree with it, down to the audit trail's timestamps. The lock
behaviour is asserted against a real server holding a real conflicting lock, not
argued. Compression against PostgreSQL row storage is **measured** rather than
targeted β ~109Γ (457 B/row against 4.2 B/row; the measurement suite carries the
caveats).
Missing: query-latency benchmarks on reference hardware, so the p99 targets remain
aspirational; Spark and Trino interop. Compaction and general orphan-file cleanup
[run out of band](https://hupe1980.github.io/meterstore/docs/operations/#compaction),
because `iceberg-rust` exposes neither.
## Development
Requires a Rust toolchain and, for integration tests, a running Docker daemon.
```bash
just # list all recipes
just dev # format + unit tests (no Docker)
just test # full suite
just check # everything CI runs
just cli status # run the command-line tool from source
just site # serve the documentation site
```
The integration suites are **one test binary against one PostgreSQL container**.
Cargo would otherwise compile each file under `tests/` into its own statically
linked executable β twenty-odd full copies of DataFusion, Arrow, Iceberg, sqlx and
tonic, about 9 GB, which exhausts a CI runner's disk and fails as a linker bus
error rather than as "no space left". A container per test cost seven times the
wall clock of a database per test, for the same isolation.
`datafusion`, `arrow`, `iceberg`, `parquet`, `metering`, `time`, `rust_decimal`
and `sqlx` must each appear exactly once in the dependency graph; `just deps`
fails the build otherwise. Two versions of `arrow` mean two incompatible
`RecordBatch` types, and two of `sqlx` mean two incompatible `PgPool` types β
neither fails obviously.
## License
Dual-licensed under [MIT](LICENSE-MIT) or [Apache-2.0](LICENSE-APACHE), at your
option. Part of the [mako](https://github.com/hupe1980/mako) platform.