Files
swarm-house/docs/03-data-platform.md
T
eSlider 87551fc1c1
CI & Release / Verify simulator (push) Skipped
CI & Release / Trivy scan (push) Skipped
CI & Release / Semantic Release (push) Skipped
docs: put MinIO first-class on ground, add executive overview
Treat MinIO as existing stack infrastructure (default T3 warehouse, not
on the drone), add a one-page executive map, and record placement as an
open team question rather than a rip-and-replace.
2026-07-17 12:22:21 +01:00

146 lines
8.8 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 03 — Data platform
The heart of the proposal: how raw sensor streams become a compressed, queryable, offloadable store — with **one format and one layout on every floor** of the system.
## Datasets
Three logical datasets, distinct because their consumers and sync policies differ:
| Dataset | Contents | Rate | Leaves the drone in flight? |
| --- | --- | --- | --- |
| `telemetry` | Raw sensor rows: IMU, barometer, temperature, LiDAR, RSSI, power, magnetometer, … | Up to hundreds of Hz per sensor | **Never** |
| `detections` | Object-detection events from the video service: class, confidence, bounding box, ego-position at detection time | Event-shaped, bursty | Yes — derived, shareable |
| `state` | The drone's own pose stream: position, attitude, velocity, timestamp; plus the log of every broadcast sent/received | 110 Hz effective | Yes — this *is* the shared state |
## Partitioning layout
Identity dimensions first, calendar time last, everything Hive-style `key=value`:
```
dataset=telemetry/flight=20260708T1100Z-a3f2/drone=dr-017/sensor=imu/year=2026/month=07/day=08/hour=11/
current/min_00.parquet ← hot write path (open window)
current/min_01.parquet
data.parquet ← sealed once the hour closes
dataset=detections/flight=…/drone=…/year=…/month=…/day=…/hour=…/
dataset=state/flight=…/drone=…/year=…/month=…/day=…/hour=…/
```
Design decisions and their reasons:
- **`flight` before `drone`.** The dominant warehouse query is *replay of one flight across all drones*. With `flight` on top, a replay reads one subtree. (On board, `drone=` is constant and costs nothing — but keeping it means the on-board layout is byte-identical to the warehouse layout.)
- **`sensor` before time.** Each sensor has its own schema, so sensors need separate leaves anyway; and "one sensor over a period" is the most common analytical scan. Partition pruning handles both.
- **Flight identifiers.** The canonical flight id is a **UUIDv7** (time-ordered, globally unique, embedded timestamp — sortable by design). Partition paths use a short human-readable time-sortable alias (`20260708T1100Z-a3f2`) so directory listings stay debuggable; the full UUIDv7 lives in the data and in the flight manifest.
- **Timestamps: `int64` epoch nanoseconds.** Same 8 bytes as milliseconds, native resolution of ROS 2, sufficient for IMU-class rates, valid until year 2262. One time type everywhere — no unit confusion at merge time.
- **No rounding in raw data.** Coordinates and measurements are stored exactly as the sensor emits them (float64/float32 per sensor spec). ZSTD plus Parquet delta encoding makes rounding-for-size unnecessary. Quantization (millimeters as int32) is applied **only** to broadcast payloads, where every byte of air time counts ([04 — Swarm sync](04-swarm-sync.md)).
## Hot write path: `current/` → sealed
Streaming many small rows straight into big Parquet files is impossible (Parquet is immutable); buffering a whole hour in RAM is unacceptable (crash = lost data). The proven pattern:
```mermaid
sequenceDiagram
participant S as sensor-ingest
participant W as parquet-writer
participant C as current/ (open window)
participant Z as sealer
participant P as data.parquet (sealed)
participant D as DuckDB
loop every minute block
S->>W: rows (int64 ns, typed)
W->>C: flush min_NN.parquet (small, ZSTD-fast)
end
Note over Z: window closes (hour sealed<br/>or flight ends)
Z->>C: read all min_*.parquet
Z->>P: write one compact file (ZSTD-high)
Z->>C: remove blocks after verify
D->>P: read_parquet('…/**/*.parquet') — glob covers both states
```
- Crash tolerance is one minute-block, tunable per sensor rate.
- Readers never care about the sealing state: a single `read_parquet` glob matches both `current/*.parquet` and `data.parquet`.
- The sealer runs at low priority — flight-critical software always wins the CPU.
## Storage floors (tiers)
Same format on every floor; only volume, retention, and location change:
| Floor | Where | Contents | Retention |
| --- | --- | --- | --- |
| **T0 — hot** | RAM / DuckDB in-process | Sliding window of the last minutes; what mission logic queries in flight | Minutes |
| **T1 — warm** | Drone NVMe (Parquet) | Full raw `telemetry` + `detections` + `state` of the current flight; never leaves the drone in flight | Current flight (+ quota-based headroom) |
| **T2 — shared** | Peer-derived cache on NVMe (Parquet) | Derived data only (`detections`, `state`) pulled from peers when the link allows | Current mission |
| **T3 — warehouse** | Ground **MinIO** (or equivalent) + Parquet lakehouse | Every flight of every drone; replay, analytics, model training — **primary MinIO home** | Years |
| **T4 — dev** | Engineer laptop / sim farm | Slices pulled from T3, or synthetic data from the simulator | Ephemeral |
**MinIO note:** default sketch keeps MinIO on the **ground (T3)**. An on-board MinIO sidecar is a possible exception if the team already standardises on an S3 API in Compose — derived data only, never the 5 Hz pose path. Placement is an open question ([09](09-open-questions.md), [00 — Executive](00-executive.md)).
## Flight offload: a mirror, not a migration
Because T1 and T3 share the identical layout, offload after landing is:
```bash
# per drone, at the base station
mc mirror drone-nvme/flights/ warehouse/flights/ # or rsync over the wired dock
duckdb -c "…integrity audit: row counts, time coverage, gap scan per partition…"
```
The audit compares expected vs actual coverage per partition (row counts, min/max timestamps, gap detection) before the drone's local copy is released for cleanup. Missing or corrupt partitions are re-pulled — a lazy gap-repair pass rather than a failed batch job.
## Retention and quotas on board
- Per-dataset NVMe quotas enforced by the sealer (oldest sealed partitions of *previous* flights evicted first; the current flight is never evicted).
- If ingest pressure exceeds the write budget, degradation is explicit and ordered: reduce optional high-rate sensors' sampling first, never drop `state` or `detections`.
## Query examples (DuckDB)
```sql
-- All detections of one flight, fused across the fleet (runs on T3)
SELECT drone, ts_ns, class, confidence, pos_x, pos_y, pos_z
FROM read_parquet('warehouse/dataset=detections/flight=20260708T1100Z-a3f2/**/*.parquet')
ORDER BY ts_ns;
-- IMU of one drone for one minute (runs identically on T1 or T3)
SELECT *
FROM read_parquet('…/dataset=telemetry/flight=…/drone=dr-017/sensor=imu/**/*.parquet')
WHERE ts_ns BETWEEN 1783941000000000000 AND 1783941060000000000;
-- Downsample for a dashboard: 1 Hz means over a 400 Hz stream
SELECT ts_ns // 1000000000 AS second, avg(accel_x), avg(accel_y), avg(accel_z)
FROM read_parquet('…/sensor=imu/**/*.parquet')
GROUP BY second ORDER BY second;
```
## Spatial queries: DuckDB spatial extension
Positions and detections are inherently spatial, and the [DuckDB spatial extension](https://github.com/duckdb/duckdb-spatial) makes them queryable as geometry without leaving the engine (`INSTALL spatial; LOAD spatial;`). Its current lack of spherical-coordinate support is irrelevant here: the swarm operates in a **local planar mission frame in meters**, which is exactly the geometry model the extension handles best.
```sql
LOAD spatial;
-- Minimum separation between two drones over a flight (replay safety audit)
SELECT min(ST_Distance(ST_Point(a.pos_x, a.pos_y), ST_Point(b.pos_x, b.pos_y))) AS min_sep_m
FROM read_parquet('…/dataset=state/flight=…/drone=dr-01/**/*.parquet') a
JOIN read_parquet('…/dataset=state/flight=…/drone=dr-02/**/*.parquet') b
ON a.ts_ns // 200000000 = b.ts_ns // 200000000; -- align to 200 ms buckets
-- Detections inside the mission perimeter polygon
SELECT cls, count(*)
FROM read_parquet('…/dataset=detections/flight=…/**/*.parquet')
WHERE ST_Within(ST_Point(obj_x, obj_y), ST_GeomFromText('POLYGON((…))'))
GROUP BY cls;
-- Trajectory as a geometry (export for replay overlays)
SELECT drone, ST_MakeLine(list(ST_Point(pos_x, pos_y) ORDER BY ts_ns)) AS path
FROM read_parquet('…/dataset=state/flight=…/**/*.parquet')
GROUP BY drone;
```
On the ground this powers replay overlays, coverage analysis (did ubiquitous sensing actually cover the area?), and separation audits; on board it is available for cheap geo-predicates (e.g. "detections within R meters of me") with no extra service.
## Schema management
- Every dataset schema is versioned in-repo (`schemas/` — SemVer, one file per dataset version) and referenced by the fleet release manifest ([11](11-cicd-delivery.md)).
- Parquet files carry the schema version in file metadata; readers select the matching deserializer.
- Additive evolution only within a MAJOR version (new nullable columns); breaking changes bump MAJOR and require a manifest release.