Stream Data Archives
Ein Stream Data Archive ist die Einheit für Konfiguration und Speicherung von Zeitreihendaten in OctoMesh. Jedes Archiv ist eine versionierte, typisierte CrateDB-Tabelle pro Tenant, die eine kuratierte Menge von Attributpfaden eines Construction-Kit-Typs erfasst. Archive haben einen strengen Lebenszyklus (Created → Activated → Disabled / Failed), ein nach der Aktivierung unveränderliches Schema und ein dreistufiges Aktivierungs-Gate (Instanz → Tenant → Archiv), das bestimmt, ob die Datenebene geöffnet ist.
Archive sind Runtime-Entitäten – Instanzen des CK-Modells System.StreamData – und können über das Refinery Studio, GraphQL, REST oder die octo-cli erstellt, aktiviert, deaktiviert, erneut versucht und gelöscht werden.
Archiv-Zeilen können ebenfalls gesichert werden. Verwenden Sie
ExportArchiveData /
ImportArchiveData,
um die Zeilen eines einzelnen Archivs zu verschieben, oder schließen Sie alle Archive eines Tenants in ein vollständiges Tenant-Backup ein mit
dump --include-archive-data – siehe Repository Backup & Restore.
Was ein Archiv ist
Ein Archiv bündelt drei Belange in einer Entität:
- Ziel-CK-Typ – die Quelle der Wahrheit für die Datenform. Vererbung wird unterstützt: Ein auf Typ
Adefiniertes Archiv gilt für alle abgeleiteten Typen; eine CrateDB-Tabelle pro konkretem Typ. - Spaltenliste – die zur Ingestionszeit erfassten Attributpfade, abgebildet auf typisierte CrateDB-Spalten. Erforderlich vs. optional und indiziert vs. nicht indiziert wird pro Spalte konfiguriert.
- Status – der Lebenszyklus-Zustand. Inserts und Abfragen werden nur akzeptiert, wenn der Status
Activatedist.
Während ein CK-Modell die Form der Daten definiert, definiert ein Archiv, welcher Ausschnitt dieser Daten als Zeitreihe persistiert wird. Viele Archive können denselben CK-Typ adressieren – unterschiedliche Aufbewahrung, unterschiedliche Spalten, unterschiedliches Downsampling.
Eigenschaften
| Eigenschaft | Beschreibung |
|---|---|
| Pro Tenant | Jeder Tenant hat sein eigenes CrateDB-Schema. Keine tenantübergreifenden Daten. |
| Typisiert | Jeder erfasste Attributpfad wird zu einer echten CrateDB-Spalte mit deklariertem Typ. |
| Unveränderliches Schema | Nach der Aktivierung können TargetCkTypeId und Columns nicht mehr geändert werden. Neue Version = neues Archiv. |
| Multi-Archiv | Ein CK-Typ kann von beliebig vielen Archiven gleichzeitig erfasst werden. |
| Status-gesteuert | Inserts und Abfragen werden nur akzeptiert, solange der Status des Archivs Activated ist. |
| Soft-Delete bewahrt Daten | Das Deaktivieren eines Archivs blockiert Lese- und Schreibzugriffe, bewahrt aber die zugrunde liegende CrateDB-Tabelle. |
Archiv-Subtypen
Der abstrakte Basistyp Archive hat drei konkrete Subtypen (CK-Modell System.StreamData, ≥ 1.4.0):
| Subtyp | Zeitachse | Ingestionspfad | Anwendungsfall |
|---|---|---|---|
RawArchive | Ein timestamp pro Zeile | Externe Producer (SaveStreamDataInArchive@1, IStreamDataRepository.InsertAsync) | Momentanmessungen (Sensormesswerte, Maschinentelemetrie, Ereignisse) |
RollupArchive | Gebucketet [bucketStart, bucketEnd) | Nur der System-Rollup-Orchestrator – keine externen Inserts | Abgeleitete gebucketete Aggregationen eines anderen Archivs (raw oder Rollup-of-Rollup) |
TimeRangeArchive | Expliziter [from, to)-Bereich pro Zeile | Externe Producer (SaveTimeRangeStreamDataInArchive@1, IStreamDataRepository.InsertTimeRangeAsync) | Extern vorab aggregierte Daten (EDA-Berichte, Smart Meter, Wetter-APIs) |
RawArchive und TimeRangeArchive sind operatordefiniert und akzeptieren externe Daten. Die Spalten von RollupArchive werden serverseitig aus den konfigurierten Aggregationen generiert – direkte Änderungen werden abgelehnt.
CK-Modell
Wird nur geladen, wenn StreamData auf Instanzebene aktiviert ist. Liegt in System.StreamData (derzeit 1.4.0).
# Abstract base
ckType: Archive
isAbstract: true
attributes:
- TargetCkTypeId # CkId<CkTypeId>
- Columns # array<CkArchiveColumn>
- Status # CkArchiveStatus
- RawRetentionMs # int64, optional — partition-drop horizon
- MaxRetroactiveReachMs # int64, optional — bounded retro reach cap for automatic recompute (since 1.6.8); null = unbounded
# Concrete subtypes
ckType: RawArchive (derivedFrom: Archive)
ckType: TimeRangeArchive (derivedFrom: Archive)
attributes:
- Period # TimeSpan, optional — advisory hint
ckType: RollupArchive (derivedFrom: Archive)
attributes:
- SourceArchiveRtId # rtId of the archive this rollup reads from
- BucketSizeMs # int64
- BucketAlignment # FixedSize | CalendarDay | Iso8601Week | CalendarMonth | CalendarYear
- ReferenceTimeZone # IANA zone id (e.g. Europe/Vienna), optional — anchors calendar alignments (since 1.6.4)
- WatermarkLagMs # int64 — how far behind real time the orchestrator stays
- LastAggregatedBucketEnd # DateTime, optional
- Aggregations # array<CkRollupAggregation>
- FrozenUntil # DateTime, optional
ckEnum: CkArchiveStatus { Created, Activated, Disabled, Failed }
ckEnum: CkRollupFunction { Avg, Min, Max, Sum, Count, TimeWeightedAvg, StateDuration, First, Last }
ckEnum: BucketAlignment { FixedSize, CalendarDay, Iso8601Week, CalendarMonth, CalendarYear }
ckRecord: CkArchiveColumn
attributes:
- Path # attribute path, e.g. "sensor.reading.value" or "readings[*].value"
- Required # bool
- Indexed # bool (default true — opt-out only when justified)
ckRecord: CkRollupAggregation
attributes:
- SourcePath # attribute path on the source archive's CK type
- Function # Avg | Min | Max | Sum | Count | TimeWeightedAvg | StateDuration | First | Last
- TargetColumnName # optional override; otherwise derived from SourcePath + Function
Speicher-Layout
Benennung
| Element | Konvention |
|---|---|
| Schema | clean(tenantRtCkId.SemanticVersionedFullName), max. 63 Zeichen |
| Tabelle | clean(targetCkType.SVF) + "_" + clean(archiveRtCkId.SVF), max. 200 Zeichen |
| Spalte | camelCase aus der CK-Attribut-ID; verschachtelte Pfade camelCase verkettet über ColumnNameMapper |
clean() entfernt Nicht-Alphanumerische. Wenn ein generierter Name das Längenlimit überschreitet, kürzt der Helper und hängt ein deterministisches SHA-256-Suffix an ({truncated}_{hash16}). Die Kürzung wird auf Information-Ebene mit beiden Namen zur Nachvollziehbarkeit protokolliert.
Standardspalten (jede Archivtabelle)
| Spalte | Typ | Hinweise |
|---|---|---|
rtId | TEXT | Instanz-ID der Entität |
timestamp | TIMESTAMP WITH TIME ZONE | Zeitstempel des Datenpunkts (UTC) |
ckTypeId | TEXT | Konkreter Typ (relevant bei Vererbung) |
rtCreationDateTime | TIMESTAMP WITH TIME ZONE | DEFAULT CURRENT_TIMESTAMP |
rtChangedDateTime | TIMESTAMP WITH TIME ZONE | Bei Konflikt aktualisiert |
rtWellKnownName | TEXT | Optional |
Primärschlüssel: (timestamp, rtId, ckTypeId).
TimeRangeArchive ersetzt timestamp durch bucketStart + bucketEnd; der PK ist (bucketStart, bucketEnd, rtId, ckTypeId).
Pfad-zu-Spalte-Zuordnung
| Pfadform | CrateDB-Spaltentyp | Nullbarkeit |
|---|---|---|
Skalar (voltage) | abgeleitet vom primitiven Typ des CK-Attributs | NOT NULL wenn required |
Record (sensor) | OBJECT(STRICT) mit Unterfeldern aus CkRecord | wie oben |
Array von Skalaren (readings[*].value) | ARRAY(<scalar type>) | NOT NULL wenn required |
Array von Records (readings[*]) | ARRAY(OBJECT(STRICT)) | wie oben |
Jede generierte Spalte wird durch die Standardregeln von CrateDB indiziert, es sei denn CkArchiveColumn.Indexed = false, in welchem Fall das DDL INDEX OFF ausgibt. Deaktivieren Sie Indizes nur für Spalten, die gelesen, aber niemals gefiltert oder aggregiert werden.
Upsert-Semantik
Inserts laufen als mehrzeilige INSERT … ON CONFLICT DO UPDATE. Erforderliche Spalten werden immer vom eingehenden Wert überschrieben (der app-seitige Validator garantiert, dass sie vorhanden sind); optionale Spalten werden über COALESCE(EXCLUDED.x, x) zusammengeführt – mehrere Quellen können zum selben (timestamp, rtId, ckTypeId) beitragen, ohne sich gegenseitig zu überschreiben.
rtChangedDateTime wird bei jedem Konflikt hochgesetzt; rtCreationDateTime wird niemals verändert. Idempotente Re-Inserts sind sicher.
Zeitsemantik
- Die kanonische Zeitzone ist UTC. Eingehende Zeitstempel werden normalisiert; naive
DateTime-Werte werden als UTC interpretiert. - API-Antworten serialisieren als ISO-8601 mit explizitem
Z. - Verspätet eintreffende Daten werden über den Standard-Upsert-Pfad akzeptiert.
- Zukünftige Zeitstempel werden nicht abgelehnt – die Drifterkennung gehört in die Pipeline- / Monitoring-Ebene.
Lebenszyklus
| Status | Tabellenzustand | Inserts / Abfragen | Hinweise |
|---|---|---|---|
Created | Nicht provisioniert | Abgelehnt | Vollständig editierbar. Keine CrateDB-Ressourcen allokiert. |
Activated | Provisioniert, eingefroren | Akzeptiert | Schema unveränderlich. Validierung auf dem Update-Pfad. |
Disabled | Provisioniert, eingefroren | Abgelehnt | Daten bewahrt. Zum Fortsetzen erneut aktivieren. |
Failed | Kann partiell sein | Abgelehnt | Aktivierungs-DDL fehlgeschlagen. retryActivate ist idempotent. |
Pfadvalidierung läuft bei jedem Übergang, der in Activated endet (Created/Disabled/Failed → Activated): Jeder CkArchiveColumn.Path wird gegen das aktuelle CK-Modell erneut validiert. Fehler werfen ArchivePathInvalidException, bevor irgendein DDL läuft.
Löschung ist destruktiv: Die CrateDB-Tabelle wird verworfen (idempotentes DROP TABLE IF EXISTS) und die Entität wird soft-gelöscht (rtState = Archived). Historische Daten gehen verloren. Löschung ist aus jedem Status erlaubt – außer bei einem RawArchive mit aktiven Rollups, das mit RollupSourceInUseException abgelehnt wird.
Aktivierungsebenen
Jede Operation auf der Datenebene wird durch eine dreistufige, UND-gekoppelte Prüfung gesteuert:
Instance enabled? ── no ──▶ DI does not register CrateDB stack; tenant enable throws.
yes
▼
Tenant enabled? ── no ──▶ Archive activate throws.
yes
▼
Archive Activated? ── no ──▶ Insert / query throws ArchiveNotActivatedException.
yes
▼
OK
| Ebene | Quelle |
|---|---|
| Instanz | StreamData:Enabled in appsettings (Standard false) |
| Tenant | StreamDataGlobalSettings.IsEnabled (Flag pro Tenant) |
| Archiv | CkArchive.Status == Activated |
Das Studio liest beide Flags vom aggregierten Tenant-Features-Endpunkt (GET {tenantId}/v1/features/status, dessen streamData-Teil { instanceEnabled, tenantEnabled } trägt); wenn instanceEnabled = false ist, wird die Archiv-Navigation ausgeblendet. Wenn instanceEnabled = true, aber tenantEnabled = false ist, sieht der Benutzer einen leeren Zustand mit einer Aktion „Enable for tenant".
Autorisierung
Drei Rollen, im Identity Service neben AdminPanelManagement registriert:
| Rolle | Erlaubte Operationen |
|---|---|
StreamDataAdmin | Vollständiger Lebenszyklus: Archive erstellen / aktualisieren / löschen; aktivieren / deaktivieren / einschalten / erneut versuchen; Tenant aktivieren / deaktivieren; Rollup-Lebenszyklus (Freeze / Unfreeze / Rewind). |
StreamDataWriter | In aktivierte Archive einfügen. Metadaten lesen (Archive auflisten / abrufen). Kein Lebenszyklus, kein CRUD. |
StreamDataReader | Zeitreihenabfragen (simple / aggregation / grouped / downsampling). Metadaten lesen. Kein Insert, kein Lebenszyklus. |
AdminPanelManagement impliziert alle drei. Die Durchsetzung erfolgt serverseitig; das Studio steuert Buttons nur über Route Guards. Ein Benutzer ohne Schreibrechte, der eine Aktion anklickt, erhält eine Benachrichtigung aus der vom Server zurückgegebenen ArchiveNotActivatedException / Forbidden.
Rollup Archives
Ein Rollup-Archiv konsumiert ein anderes Archiv und schreibt eine Zeile pro [bucketStart, bucketEnd)-Fenster. Buckets werden durch BucketSizeMs dimensioniert und durch BucketAlignment ausgerichtet:
| Alignment | Bucket-Grenze |
|---|---|
FixedSize | LastAggregatedBucketEnd + N · BucketSizeMs (offset-invariant; ignoriert die Zone) |
CalendarDay | Lokale Mitternacht in ReferenceTimeZone |
Iso8601Week | Montag 00:00:00 Lokalzeit in ReferenceTimeZone |
CalendarMonth | Erster Tag des Monats 00:00:00 Lokalzeit in ReferenceTimeZone |
CalendarYear | 1. Januar 00:00:00 Lokalzeit in ReferenceTimeZone |
BucketSizeMs ist für die Kalender-Alignments nur informativ (Monatslängen variieren). Der Orchestrator bleibt WatermarkLagMs hinter der Echtzeit zurück, bevor er einen Bucket schließt, sodass verspätet eintreffende Rohdaten trotzdem im richtigen Fenster landen.
Bucket-Intervall vs. Quellgranularität
Wenn das Quellarchiv eine Granularität deklariert (das Period eines TimeRangeArchive oder das BucketSizeMs eines übergeordneten Rollups), muss das BucketSizeMs des Rollups ≥ der Quellgranularität und ein ganzzahliges Vielfaches davon sein. Die Aktivierung lehnt einen feineren oder nicht ausgerichteten Bucket mit RollupBucketIntervalException ab – z. B. werden über einer 15-Minuten-Quelle 15 min / 30 min / 1 h / 1 d akzeptiert, während 5 min (feiner) und 20 min (kein ganzes Vielfaches) abgelehnt werden. Die Beschränkung verhindert, dass eine einzelne Quellzeile über zwei Ausgabe-Buckets aufgeteilt wird, was Sum/Count stillschweigend verfälschen und Avg/Min/Max verzerren würde. Rohquellen, die keine deklarierte Granularität aufweisen, sind unbeschränkt. Die Prüfung liegt in RollupValidator.ValidateForActivation, sodass ein falsch konfigurierter Rollup bereits bei Activate schnell fehlschlägt, statt zur Abfragezeit falsche Aggregate zu produzieren.
Referenzzeitzone (Kalender-Alignments)
Die Kalender-Alignments rasten jeden Bucket auf eine lokale Kalendergrenze in ReferenceTimeZone ein – eine IANA-Zonen-ID wie Europe/Vienna – sodass die Buckets DST-korrekt sind: Ein CalendarDay-Bucket läuft von lokaler Mitternacht zu lokaler Mitternacht, was unter CEST (Sommer) 22:00Z → 22:00Z und unter CET (Winter) 23:00Z → 23:00Z ist. Eine leere Zone fällt auf UTC-Grenzen zurück. FixedSize ignoriert die Zone – ein Ganze-Stunden-Fenster ist für eine Ganze-Stunden-Zone offset-invariant, sodass ein stündlicher Rollup keine Zone benötigt. Verfügbar seit System.StreamData 1.6.4.
BucketAlignment und ReferenceTimeZone sind eingefroren, sobald der Rollup Activated ist. Das Umstellen eines aktiven Rollups von FixedSize/UTC auf ein Kalender-Alignment in einer echten Zone (z. B. das Verschieben eines täglichen Rollups auf CalendarDay + Europe/Vienna) erfordert das Löschen und Neuanlegen des Rollups, gefolgt von erneuter Aktivierung und Backfill.
Aggregationsfunktionen
CkRollupFunction bildet auf ein oder zwei Speicherspalten ab:
| Funktion | Gespeicherte Spalten | Formel zur Lesezeit |
|---|---|---|
Sum | {base}_sum | unverändert |
Min | {base}_min | unverändert |
Max | {base}_max | unverändert |
Count | {base}_count (BIGINT) | unverändert |
Avg | {base}_sum, {base}_count (beide gespeichert) | sum / count – hält verkettete Rollups numerisch korrekt |
First | {base}_first | Wert bei der frühesten Beobachtung im Bucket (numerische Quellspalten) |
Last | {base}_last | Wert bei der spätesten Beobachtung im Bucket (numerische Quellspalten) |
Die Columns[]-Projektion eines Rollups wird über RollupColumnGenerator.Generate aus Aggregations[] generiert; direkte Änderungen an Columns werden vom Validierungs-Hook abgelehnt.
Jede Deklaration in Aggregations[] ist unabhängig, sodass ein Rollup-Archiv viele Aggregationen über verschiedene Attribute tragen kann (z. B. Sum von energy und Max von power), plus mehrere Funktionen desselben Attributs (z. B. Min und Max von power). Alle deklarierten Aggregate werden in einem einzigen Durchlauf über jeden Quell-Bucket erzeugt und jedes wird als eigene Serie bereitgestellt. Um einem bereits aktivierten Rollup eine Aggregation hinzuzufügen, erstellen Sie einen neuen Rollup mit der vollständigen Aggregationsmenge und füllen ihn mit rewindRollupWatermark per Backfill – Aggregations[] ist nach der Aktivierung unveränderlich.
Operator-Steuerungen
| Operation | Mutation | Hinweise |
|---|---|---|
| Create | createRollupArchive(input) | Der Server leitet TargetCkTypeId und Columns aus dem Quellarchiv + Aggregationen ab. |
| Bereich einfrieren | freezeRollupArchive(rtId, until) | Monoton – das neue until muss ≥ dem aktuellen FrozenUntil sein. Der Orchestrator hört auf, Buckets zu produzieren, deren bucketEnd in den eingefrorenen Bereich fällt. |
| Unfreeze | unfreezeRollupArchive(rtId, acceptGaps) | Idempotent. acceptGaps wird für die Auditierung erfasst, aber der Lückenerkennungs-Guard ist ein Follow-up. |
| Watermark zurückspulen | rewindRollupWatermark(rtId, toBucketEnd) | Aggregiert ab der angegebenen Bucket-Grenze vorwärts neu. Destruktiv: Zuvor committete Zeilen im zurückgespulten Bereich sind vorübergehend inkonsistent, bis der Orchestrator aufholt. |
| Bereich neu berechnen | recomputeArchive(rtId, from, to, rtIdScope?) | Optimistische, ausfallfreie Neuberechnung von [from, to), optional auf eine einzelne Entität eingegrenzt. Leser sehen durchgängig einen konsistenten Snapshot. Siehe Rollup Recompute unten. |
Ein RawArchive, auf das ein aktiver Rollup verweist, kann nicht gelöscht werden – löschen Sie zuerst den Rollup.
Rollup Recompute
Wenn sich Rohdaten nachträglich ändern – verspätete Werte, Korrekturen, erneute Ingestionen – sind die aus diesem Fenster abgeleiteten Rollup-Zeilen veraltet. Eine Recompute aggregiert die betroffenen Fenster neu. Die harte Anforderung ist, dass Konsumenten (Meshboards, Pipelines, die API) während eine Recompute läuft weiterarbeiten: kein 500, keine abgebrochene Pipeline, und jeder Lesevorgang liefert einen konsistenten Snapshot – entweder die vorherigen Werte oder die neuen, niemals eine halb geschriebene Mischung.
Für die ingenieursseitigen Interna – wie ein retroaktiver Schreibvorgang erkannt wird (RetroactiveWriteDetector, das Kriterium der konsumierten Watermark), das Dirty-Windows-Ledger, die Propagation über den Abhängigkeitsgraphen, der periodische Orchestrator, der optimistische atomare Swap mit einem End-to-End-Sequenzdiagramm und einem Vergleich aller vier (Re-)Compute-Pfade – siehe Rollup Recompute Internals.
Trigger
| Trigger | Oberfläche | Verhalten |
|---|---|---|
| Manuell | GraphQL-Mutation recomputeArchive · POST /archives/{rollupRtId}/recompute · octo-cli RecomputeArchive | Ein Operator erzwingt eine Recompute über [from, to), optional auf eine Entität eingegrenzt (rtIdScope). Erfordert StreamDataAdmin. |
| Periodisch | RecomputeOrchestratorHostedService (Tick pro Tenant) | Leert das Dirty-Windows-Ledger und berechnet nur die tatsächlich als dirty markierten Fenster neu – nicht das gesamte Archiv. |
| Kettenpropagation | intern | Nachdem die Recompute eines Rollups committet wurde, werden seine eigenen Abhängigen als dirty markiert, sodass eine mehrstufige Rollup-of-Rollup-Kette von oben nach unten weiterrollt. |
Der periodische Trigger ist der automatische Pfad – er leert das Dirty-Windows-Ledger, das der Detector bei jedem retroaktiven Schreibvorgang füllt, sodass kein Operator einen Schwellenwert von Hand zurücksetzen muss. Ein einzelner sehr verspäteter Schreibvorgang kann eine automatische Recompute nur so weit zurückziehen wie die pro Archiv geltende beschränkte Retro-Reichweite-Obergrenze (Archive.MaxRetroactiveReachMs, seit 1.6.8; null = unbeschränkt), kombiniert mit der Host-Obergrenze StreamData:Recompute:MaxRetroactiveReachHardLimitMs. Manuelles recomputeArchive und rewindRollupWatermark bleiben unbeschränkt – die bewusste Notausstiegsluke für eine tiefe Korrektur (AB#4196).
Optimistischer atomarer Swap (wie Konsistenz garantiert wird)
CrateDB hat keine Mehr-Statement-Transaktion, daher verwendet der Swap einen Generation-Pointer pro Fenster:
- Jede Rollup-Tabelle trägt eine
generation-Spalte (BIGINT NOT NULL DEFAULT 0), die Teil des Primärschlüssels ist. Die normale Vorwärtsaggregation schreibt immergeneration = 0. - Eine Recompute aggregiert den Bereich in eine Staging-Tabelle pro Archiv und kopiert diese Zeilen dann in die Live-Tabelle, gestempelt mit der nächsten Generation
N+1. Die vorherige Generation bleibt an Ort und Stelle und sichtbar. - Eine winzige Side-Table pro Rollup (
archive_<rtId>__genmap) hält die aktive Generation pro Bereich. Das Umschalten dieses Pointers aufN+1ist ein Einzelzeilen-Schreibvorgang – der atomare Commit-Punkt. - Der Lesepfad injiziert
generation = CASE WHEN <range> THEN <activeGen> … ELSE 0 END, sodass jede Abfrage genau eine konsistente Generation pro Fenster auflöst. - Nach dem Umschalten löscht ein Hintergrund-Sweep die nun überholten Zeilen der alten Generation.
Fehlersemantik: Ein Job, der vor dem Umschalten abstürzt, hinterlässt die Staging-Tabelle verwaist (beim nächsten Lauf über den Namen zurückgewonnen) und die Live-Tabelle + den Pointer unberührt – Leser sehen weiterhin die vorherigen Werte. Ein Absturz nach dem Umschalten, aber vor dem Sweep, hinterlässt nur tote Zeilen, die der nächste Sweep zurückgewinnt. Es ist niemals ein partieller Zustand beobachtbar.
Scoping
rtIdScope beschränkt die gesamte Pipeline – Staging-Aggregation, den Generation-Pointer und den Sweep nach dem Umschalten – auf eine einzelne Entität (Zählpunkt / Stream). Ein verspäteter Wert für einen Stream berechnet daher nur die Fenster dieses Streams neu und berührt niemals unbeteiligte Entitäten, die sich denselben Bucket teilen. Das Weglassen von rtIdScope berechnet den gesamten Bereich neu.
Coalescing
Überlappende Trigger für dasselbe Archiv verschmelzen (coalesce): Ihre Bereiche werden zu einem aktiven Job zusammengeführt, und der überholte Trigger wird mit dem Zustand Coalesced in der Job-Historie erfasst. Es gibt zu jedem Zeitpunkt höchstens eine aktive Recompute pro Archiv.
Interaktion mit Rewind
rewindRollupWatermark aggregiert einen Bereich bei generation = 0 vorwärts neu. Wäre dieser Bereich zuvor neu berechnet worden, würde der Aktive-Generation-Pointer die Leser andernfalls weiterhin an die veraltete neu berechnete Generation pinnen. Das Rewind gleicht daher die Generation-Map ab: Es entfernt die Pointer, die über die Rewind-Grenze hinausreichen, und löscht die verwaisten generation > 0-Zeilen dort, sodass die Vorwärts-Neuaggregation wieder maßgeblich wird. Dies ist ein No-op, wenn Streamdaten deaktiviert sind oder nichts neu berechnet wurde.
Job-Historie & Observability
Jede Recompute wird als Job-Snapshot erfasst, abfragbar neueste-zuerst (auf 50 begrenzt):
- GraphQL –
runtime.streamData.recomputeJobsFor(rtId)liefert[RecomputeJobInfo!]!. - REST –
GET /archives/{archiveRtId}/recompute-jobs. - octo-cli –
ListRecomputeJobs.
RecomputeJobInfo trägt state (Running / Completed / Failed / Coalesced), rowsProcessed, windowsProcessed, startedAt, finishedAt, durationMs und errorReason (bei Fehler gesetzt). Zusätzlich stellt das pro Rollup geltende RollupArchiveInfo (über rollupsFor) Live-Health-Felder bereit: recomputeInProgress, lastRecomputeStartedAt, lastRecomputeSuccessAt, lastRecomputeFailureAt, lastRecomputeFailureReason, dirtyWindowsPending und pendingRecomputeRanges. Fehler werden außerdem über den Archiv-Audit-Trail ausgegeben.
Time-Range Archives
Für extern vorab aggregierte Daten mit expliziten [from, to)-Fenstern. Auf diesem Pfad gibt es kein Quellarchiv und keinen Orchestrator; Producer pushen Zeilen direkt:
- GraphQL-Mutation:
createTimeRangeArchive(input)– derselbe Admin-stufige Rollen-Guard wie bei Rollup-Create. - Pipeline-Node:
SaveTimeRangeStreamDataInArchive@1. - REST:
POST /api/v1/streamData/archives/{archiveRtId}/insertTimeRange.
Period auf TimeRangeArchive ist nur beratend – es lässt das Studio einen Standard-Zeitraum rendern, wird aber zur Insert-Zeit nicht durchgesetzt.
Computed Columns
Eine Computed Column auf einem Raw- oder Time-Range-Archiv wird durch eine Formel aus einer oder mehreren anderen Spalten derselben Zeile abgeleitet, statt vom Producer geliefert zu werden – z. B. powerFactor = activePower / apparentPower. Der Wert wird als echte, typisierte CrateDB-Spalte materialisiert, sodass er von Rollups nativ aggregierbar und für direkte SQL-Konsumenten (Grafana) sichtbar ist.
- Wo er ausgewertet wird – serverseitig im CrateDB-Repository auf jedem Ingest-Pfad (Pipeline-Node, REST, SDK), sodass der Wert unabhängig vom Producer identisch ist. Die Formel verwendet den geteilten mXparser-Dialekt (siehe die Referenz Formula Expressions unter Communication).
- Ergebnistyp – explizit deklariert und aus dem
double, das mXparser berechnet, zurückgecastet:Boolean,Int,Int64,DoubleoderDateTime(Ticks).Stringund nicht-skalare Typen werden nicht unterstützt. Ein Formel- / NaN- / NULL-Operand-Fehler speichertNULLfür diese Zelle; die Ingestion schlägt nie fehl. - CK-Modell –
CkArchiveColumnträgtName+Formula+ResultType+ den Engine-verwaltetenComputedState(seitSystem.StreamData1.5.0) sowieComputedVersion/PendingFormula(seit 1.6.1 bzw. 1.6.2). Eine Spalte mit einerFormulaist berechnet; eine mit einemPathwird ingestiert. Die beiden schließen sich gegenseitig aus.
Lebenszyklus bei aktivem Archiv (optimistisch / atomar)
Das Hinzufügen, Neu-Formulieren oder Entfernen einer Computed Column funktioniert auf einem Activated-Archiv, ohne die Ingestion zu blockieren oder Leser zu brechen, indem es die optimistische Semantik von Rollup Recompute wiederverwendet:
- Add – die Spalte wird
Pendingpersistiert, die physische Spalte wird hinzugefügt (ALTER TABLE ADD COLUMN), ein Backfill befüllt sie über die bestehenden Zeilen, während sie für Leser verborgen bleibt, dann schaltet sie atomar aufActive. Ein Backfill-Fehler hinterlässt sieFailedund den vorherigen Archivzustand intakt. - Formel ändern – die Spalte liefert weiterhin ihre vorherigen Werte (Version N), während die neue Formel in eine versionierte physische Spalte
{base}__v{N+1}per Backfill gefüllt wird; die Ingestion dual-writet während des Backfills in beide; dann tauscht ein einzelner Store-Schreibvorgang den Versions-Pointer atomar aus. Ein Fehler kehrt zur vorherigen Formel zurück. - Remove – aus der logischen Projektion entfernt; die physische Spalte bleibt als harmloser nullbarer Waise zurück, den der Lesepfad nicht mehr projiziert.
Einschränkung (Formeländerung). Da eine Formeländerung den physischen Namen der Spalte neu versioniert, wird das Ändern der Formel einer Computed Column, die eine andere Computed Column referenziert, abgelehnt (ComputedColumnInvalidException) – zeigen Sie zuerst den Abhängigen um oder entfernen Sie ihn. Direct-SQL- / Grafana-Konsumenten, die an den physischen Namen binden, müssen der Umbenennung ebenfalls folgen.
Rollups über Computed Columns
Ein CkRollupAggregation.SourcePath kann über ihren Name eine Quell-Computed-Column adressieren – der Rollup aggregiert sie nativ wie jede ingestierte Spalte (single- und multi-aggregate). Ein RollupArchive kann außerdem eigene Computed Columns über seine aggregierten Ausgabespalten tragen.
API-Oberfläche
GraphQL – Lebenszyklus-Mutationen
Unter mutation.streamData:
| Mutation | Liefert | Beschreibung |
|---|---|---|
activateArchive(rtId) | ArchiveTransitionResult | Provisioniert die CrateDB-Tabelle, wechselt zu Activated. |
disableArchive(rtId) | ArchiveTransitionResult | Activated → Disabled. Tabelle bleibt erhalten. |
enableArchive(rtId) | ArchiveTransitionResult | Disabled → Activated. Validiert Spaltenpfade erneut. |
retryArchiveActivation(rtId) | ArchiveTransitionResult | Failed → Activated. Idempotent (CREATE TABLE IF NOT EXISTS). |
deleteArchive(rtId) | Boolean | Verwirft die CrateDB-Tabelle und soft-löscht die Entität. Abgelehnt, wenn aktive Rollups darauf verweisen. |
createRollupArchive(input) | OctoObjectId | Server-abgeleitete TargetCkTypeId + Columns. |
createTimeRangeArchive(input) | OctoObjectId | Operatordefinierter Zieltyp und operatordefinierte Spalten. |
freezeRollupArchive(rtId, until) | ArchiveTransitionResult | Monotones FrozenUntil setzen. |
unfreezeRollupArchive(rtId, acceptGaps) | ArchiveTransitionResult | Setzt FrozenUntil zurück. |
rewindRollupWatermark(rtId, toBucketEnd) | ArchiveTransitionResult | Aggregiert ab der angegebenen Grenze neu. |
recomputeArchive(rtId, from, to, rtIdScope?) | RecomputeJobInfo | Löst eine optimistische Recompute von [from, to) aus / verschmilzt sie, optional pro Entität. Liefert den Job-Snapshot. |
addComputedColumn(rtId, name, formula, resultType, indexed?) | ArchiveTransitionResult | Fügt eine Computed Column hinzu und füllt sie per Backfill (optimistisch / atomar). |
updateComputedColumnFormula(rtId, name, formula) | ArchiveTransitionResult | Formuliert eine Computed Column neu (versionierter Backfill + atomarer Swap). Abgelehnt, wenn eine andere Computed Column sie referenziert. |
removeComputedColumn(rtId, name) | ArchiveTransitionResult | Entfernt eine Computed Column (physische Spalte bleibt als Waise zurück). |
Die passende Leseabfrage ist runtime.streamData.recomputeJobsFor(rtId): [RecomputeJobInfo!]! (neueste Jobs, neueste zuerst, auf 50 begrenzt).
Alle Mutationen erfordern StreamDataAdmin. Domänenausnahmen (§Exceptions unten) treten als stabile GraphQL-Fehlercodes über error.extensions.OctoDetails zutage.
CRUD auf Archiven (List / Get / Update) wird automatisch aus dem CK-Modell generiert und liegt unter runtime.systemStreamDataArchives / …RawArchives / …RollupArchives / …TimeRangeArchives. Der Validierungs-Hook erzwingt die Unveränderlichkeit nach der Aktivierung.
GraphQL – Pfadaufzählung
query availableArchivePaths(ckTypeId: ID!, maxDepth: Int = 5): [ArchivePathInfo!]!
type ArchivePathInfo {
path: String! # "sensor.reading.value", "readings[*].value"
primitiveType: String # null when isRecord = true
isRecord: Boolean!
isArray: Boolean!
recordTypeId: String
inheritedFromCkTypeId: String # null if attribute belongs to ckTypeId itself
}
Der AttributePathPicker des Studios konsumiert dies für das Erstellen/Bearbeiten von Archiven. Nur für Admins.
REST – /api/v1/streamData
| Methode | Pfad | Beschreibung |
|---|---|---|
| POST | /enable / /disable | Umschalter auf Tenant-Ebene. /disable antwortet mit 409, solange noch irgendein Archiv Activated ist (AB#4255). |
| POST | /archives/{archiveRtId}/activate | Wie GraphQL activateArchive. |
| POST | /archives/{archiveRtId}/disable / /enable / /retry | Statusübergänge. |
| DELETE | /archives/{archiveRtId} | Wie GraphQL deleteArchive. |
| POST | /archives/{rollupRtId}/freeze / /unfreeze / /rewind | Nur-Rollup-Steuerungen. |
| POST | /archives/{rollupRtId}/recompute | Optimistische Recompute von [from, to) (rtIdScope optional). Liefert den Job-Snapshot. |
| GET | /archives/{archiveRtId}/recompute-jobs | Aktuelle Recompute-Jobs (neueste zuerst) zum Debuggen. |
| POST | /archives/{archiveRtId}/insertTimeRange | Bulk-Time-Range-Insert. |
| GET | /archives/{archiveRtId}/rollups | Aktive Rollups, die auf dieses Raw-Archiv verweisen. |
Der Enabled-Zustand wird über das aggregierte GET {tenantId}/v1/features/status gelesen (AB#4884, ersetzt das
frühere GET /api/v1/streamData/status): Es meldet alle vier Tenant-Fähigkeiten, die der Tenant-Delete/Detach-Guard auswertet –
streamData mit { instanceEnabled, tenantEnabled }, plus communication, reporting und aiServices
mit jeweils { tenantEnabled }. Ein Lesefehler antwortet mit 500; der Zustand wird niemals als „disabled" gemeldet, wenn er
nicht gelesen werden kann.
Pipeline-Node – SaveStreamDataInArchive@1
Ein PathNodeConfiguration-Node mit einer erforderlichen Eigenschaft ArchiveRtId : OctoObjectId. Der Node leitet jeden eingehenden StreamDataPoint über IStreamDataRepository.InsertAsync an das konfigurierte Archiv weiter. Der Gegenpart SaveTimeRangeStreamDataInArchive@1 behandelt TimeRangeArchive.
Das Service Account des Mesh Adapters hält StreamDataWriter. Das Adressieren eines inaktiven Archivs führt zu einem klaren Pipeline-Fehler (ArchiveNotActivatedException).
SDK – IStreamDataRepository
// Data plane
Task InsertAsync(OctoObjectId archiveRtId, StreamDataPoint point);
Task InsertAsync(OctoObjectId archiveRtId, IEnumerable<StreamDataPoint> points);
Task InsertTimeRangeAsync(OctoObjectId archiveRtId, IEnumerable<TimeRangeStreamDataPoint> points);
Task<StreamDataQueryResult> ExecuteQueryAsync(OctoObjectId archiveRtId, StreamDataQueryOptions options);
Task<StreamDataQueryResult> ExecuteAggregationQueryAsync(OctoObjectId archiveRtId, StreamDataAggregationQueryOptions options);
Task<StreamDataQueryResult> ExecuteGroupedAggregationQueryAsync(OctoObjectId archiveRtId, StreamDataGroupedAggregationQueryOptions options);
Task<StreamDataQueryResult> ExecuteDownsamplingQueryAsync(OctoObjectId archiveRtId, StreamDataDownsamplingQueryOptions options);
// Control plane (called by IArchiveLifecycleService)
Task EnsureArchiveCreatedAsync(CkArchive archive);
Task DeleteArchiveAsync(CkArchive archive);
Bulk-Inserts verwenden Alles-oder-nichts-Vorvalidierung: Jeder Punkt wird validiert, bevor irgendein SQL gesendet wird; ein Verstoß wirft RequiredAttributeMissingException mit dem Index des betreffenden Punkts – keine Zeile wird geschrieben. Ein Fehler auf DB-Ebene nach der Vorvalidierung tritt als Exception zutage.
Bulk-Insert-Semantik
- Archiv auflösen (einzelner Lookup), Statusprüfung.
- Jeden Punkt im Batch vorvalidieren – Abdeckung erforderlicher Pfade, Typkompatibilität, Pfadauflösung. Stoppt beim ersten Verstoß.
- Bei Verstoß → mit dem betreffenden Index werfen. Keine Zeile geschrieben.
- Bei Erfolg → einzelne mehrzeilige
INSERT … ON CONFLICT …-Ausführung.
Pipeline-Aufrufer sehen Batch-Erfolg oder vollständigen Fehlschlag; Szenarien mit partiellem Commit beschränken sich auf transiente DB-Fehler (operativ, siehe Reconciliation unten).
Aufbewahrung
Die Aufbewahrung ist partitionsbasiert und pro Archiv (Work in Progress – der Entwurf ist hier festgehalten; einige Attribute sind im Create-Formular noch nicht verfügbar).
- DDL: Archivtabellen sind
PARTITIONED BY DATE_TRUNC('day', timestamp). Drop-Partition ist auf CrateDB nur Metadaten. RawRetentionMsauf dem BasistypArchivedefiniert den Aufbewahrungshorizont.null= für immer aufbewahren.- Der
ArchiveRetentionScheduler(IHostedService, stündlich) verwirft pro Archiv und pro Tenant Partitionen, die älter alsnow() - RawRetentionMssind. Fehler sind isoliert; ein fehlerhaftes Archiv blockiert die übrigen nicht.
Für Aufbewahrung im Downsampling-Stil (Rohzeilen werden vor der Löschung in ein gröberes Archiv gefaltet) verwenden Sie ein RollupArchive, das auf das Quellarchiv verweist – der Rollup ist der Langzeitspeicher, während die Aufbewahrung des Raw-Archivs aggressiv sein kann.
Reconciliation
CkArchive-Entitäten leben in MongoDB; Archivtabellen leben in CrateDB. Es gibt keine gemeinsame Transaktion. Der Lifecycle-Service hält die beiden ohne 2PC abgleichbar:
- Alles DDL ist idempotent (
IF EXISTS/IF NOT EXISTS). - Die Operationsreihenfolge ist Crate-zuerst, Mongo-zuletzt:
ActivateführtCREATE TABLEaus, dann schaltet esStatusin Mongo um.DeleteführtDROP TABLEaus, dann soft-löscht es die Entität. Mongo ist die Quelle der Wahrheit. - Jeder Aufruf auf der Datenebene prüft
Status == Activated, bevor Crate berührt wird. Solange Schritt 1 gelaufen ist, Schritt 2 aber nicht, lehnt die API mitArchiveNotActivatedExceptionab. - Beim Start zählt der
ArchiveReconcilerdieActivated-Archive pro Tenant auf und (re-)provisioniert alle fehlenden CrateDB-Tabellen. Drift in die andere Richtung (Tabelle ohne Entität) wird protokolliert, aber nicht automatisch verworfen.
Nebenläufigkeit
| Szenario | Verhalten |
|---|---|
Zwei parallele Activate auf demselben Archiv | Serialisiert über RepositoryDistributedLockService. Der zweite Aufruf beobachtet Activated und kehrt idempotent zurück. |
Disable, während Inserts laufen | Laufende Inserts werden abgeschlossen. Neue Inserts nach dem Umschalten lösen ArchiveNotActivatedException aus. |
Delete, während Inserts laufen | Wie Disable, plus das DROP TABLE kann einmal mit kurzem Backoff wiederholen, falls die Tabelle kurzzeitig gesperrt ist. |
| Activate schlägt mitten im DDL fehl | Status → Failed. EnsureArchiveCreatedAsync ist idempotent (CREATE TABLE IF NOT EXISTS), sodass ein Retry sicher ist. |
| Bulk-Insert auf deaktiviertem Archiv | Gesamter Batch beim einzelnen Status-Lookup abgelehnt. |
Exception-Katalog
Alle in Meshmakers.Octo.Runtime.Contracts.StreamData. Gemeinsame Basis: StreamDataException (erweitert RuntimeRepositoryException). Treten über GraphQL error.extensions.OctoDetails und über REST-Statuscodes zutage.
| Exception | Geworfen, wenn |
|---|---|
StreamDataException | Basis – generischer StreamData-Fehler. |
StreamDataNotEnabledException | Instanz- oder Tenant-Flag ist false bei einer Operation, die es erfordert. |
ArchiveNotFoundException | archiveRtId löst sich nicht zu einer Archive-Entität auf. |
ArchiveNotActivatedException | Archiv existiert, aber Status ≠ Activated beim Insert / bei der Abfrage. |
ArchiveSchemaImmutableException | Ein Update nach der Aktivierung versucht, schemarelevante Felder zu ändern. |
ArchivePathInvalidException | Ein CkArchiveColumn.Path kann nicht gegen TargetCkTypeId aufgelöst werden. |
ArchiveColumnTypeUnsupportedException | Ein Attributtyp kann nicht auf einen CrateDB-Spaltentyp abgebildet werden. |
RequiredAttributeMissingException | Ein Insert hat keinen Wert für einen required: true-Skalarpfad. Trägt den Index des betreffenden Punkts. |
ArchiveActivationFailedException | Die DDL-Ausführung während der Aktivierung ist fehlgeschlagen. Umschließt den zugrunde liegenden SQL-Fehler. |
InvalidArchiveStateTransitionException | Eine Lebenszyklus-Mutation wurde gegen einen inkompatiblen Status aufgerufen. |
RollupSourceInUseException | deleteArchive auf einem Raw-Archiv aufgerufen, das noch aktive Rollups hat. |
RollupSourceMissingException | createRollupArchive verweist auf eine Quell-archiveRtId, die nicht existiert oder noch nicht aktiviert ist. |
Observability
Metriken
| Metrik | Typ | Tags |
|---|---|---|
streamdata.archive.count | Gauge | tenant, status |
streamdata.archive.status_transitions | Counter | tenant, archive, from, to |
streamdata.archive.activation_duration_ms | Histogram | tenant, archive, outcome |
streamdata.insert.duration_ms | Histogram | tenant, archive, batch_size_bucket |
streamdata.insert.points | Counter | tenant, archive |
streamdata.insert.required_violations | Counter | tenant, archive |
streamdata.query.duration_ms | Histogram | tenant, archive, query_type |
streamdata.query.rows_returned | Histogram | tenant, archive, query_type |
streamdata.crate.connections.open | Gauge | (instanzweit) |
streamdata.crate.retention_partitions_dropped | Counter | tenant, archive |
streamdata.crate.rollup_rows_inserted | Counter | tenant, rollup_archive |
Traces
Activity-Sources: Meshmakers.Octo.StreamData (Top-Level-Operationen) und Meshmakers.Octo.StreamData.Crate (DB-Aufrufe). Span-Attribute: streamdata.archive.rtid, streamdata.archive.target_cktype, streamdata.tenant.
Audit-Trail
Statusübergänge werden über IEventRepository als RtEvent-Records mit source = RtEventSourcesEnum.StreamData ausgegeben. Der auslösende Benutzer ist in der Ereignisnachricht kodiert; Auto-Übergänge verwenden system:reconciliation. Die Archiv-Detailansicht des Studios stellt die Historie über den Standard-Ereignis-Abfragepfad bereit.
Ein Archiv erstellen
Über die CLI
# Tenant-level enable (required once per tenant)
octo-cli -c EnableStreamData
# Archives are runtime entities — create them through the studio, GraphQL,
# or by rt-importing a runtime-model YAML. Then:
octo-cli -c ActivateArchive -id 69fda707d47638c68edc7fea
Siehe octo-cli — ActivateArchive und verwandte Archivbefehle in der Seitenleiste für die vollständige Befehlsoberfläche.
Über das Studio
Navigieren Sie zu Repository → Stream Data Archives. Die Liste zeigt jedes Archiv des aktiven Tenants mit einem Status-Badge und Aktionen pro Zeile (Activate, Disable, Enable, Retry, Delete, Copy ID). Die Einstiegspunkte New Archive, New Rollup und New Time-Range führen durch den Create-Flow mit dem AttributePathPicker zur Spaltenauswahl.
Siehe Studio: Archives für die UI-Anleitung.
Über GraphQL
# 1. Create the archive via the auto-generated runtime mutation
mutation {
runtime {
systemStreamDataRawArchives {
create(input: {
rtWellKnownName: "EnergyMeterArchive",
targetCkTypeId: "Industry.Energy/EnergyMeter",
columns: [
{ path: "voltage", required: true, indexed: true },
{ path: "current", required: true, indexed: true },
{ path: "powerFactor", required: false, indexed: false }
]
}) { item { rtId } }
}
}
}
# 2. Activate
mutation {
streamData {
activateArchive(rtId: "65d5c447b420da3fb1238201") {
rtId
status
}
}
}
Best Practices
- Ein Archiv pro Belang. Unterschiedliche Aufbewahrung, unterschiedliche Abfragemuster, unterschiedliche Konsumenten → unterschiedliche Archive.
- Seien Sie explizit bei required. Erforderliche Skalare werden zu
NOT NULL-Spalten und werden app-seitig vor dem Insert validiert. - Indizieren Sie sparsam. Indizes kosten Speicher und Insert-Durchsatz. Standardmäßig indiziert; verzichten Sie darauf für Spalten, die gelesen, aber niemals gefiltert werden.
- Rollen Sie früh auf. Ein
RollupArchiveist über lange Horizonte weitaus günstiger abzufragen als ein Raw-Scan. - Wählen Sie das Alignment für den Anwendungsfall. Energie- / EDA-Workflows wollen Kalender-Alignments (monatlich, wöchentlich); Engineering-Telemetrie will üblicherweise
FixedSize. - Frieren Sie vor Bulk-Importen ein. Wenn Sie die Historie eines Quellarchivs per Backfill füllen, frieren Sie alle abhängigen Rollups ein, beenden Sie den Import und heben Sie das Einfrieren dann mit
acceptGaps: falseauf (oder spulen Sie die Watermark zurück), um sauber neu zu aggregieren. - Zuerst soft-löschen. Deaktivieren Sie ein Archiv, um zu bestätigen, dass kein Konsument bricht, bevor Sie löschen.
Siehe auch
- Rollups & recompute – Erkennungskriterium, Dirty-Ledger, Abhängigkeitsgraph, optimistischer atomarer Swap
- Stream Data Access (Überblick) – Abfrageoberflächen (typed / transient / persisted)
- Simple Query – rohe Zeitreihenabfrage
- Aggregation Queries – serverseitige Aggregation
- Downsampling Query – gebucketetes Downsampling
- octo-cli — EnableStreamData und verwandte Archivbefehle (Activate / Disable / Enable / Retry / Delete) – CLI-Befehle für Archivoperationen
- Studio: Archives – UI-Anleitung im Refinery Studio