Zum Hauptinhalt springen

Wie Stream Data Archives funktionieren

Diese Seite ist die architektonische Grundlage für Streamdaten: das mentale Modell und das Vokabular, das Sie benötigen, bevor Sie entweder die vollständige Referenz Stream Data Archives oder das Deep-Dive Rollups & recompute lesen. Sie erklärt, wie die Teile zusammenpassen – die Aufteilung auf zwei Stores, die drei Archivtypen, wie Zeilen gespeichert werden, wie ein Archiv seinen Lebenszyklus durchläuft und was auf dem Schreib- und Lesepfad geschieht – und übergibt dann an die Referenz- und Deep-Dive-Seiten für die erschöpfenden Details.

Wo die Details liegen

Diese Seite bleibt bewusst auf Architekturhöhe. Für die feldweise Referenz (das vollständige CK-Modell-YAML, Spaltentabellen, den Exception-Katalog, Metriken, die API-Oberfläche) lesen Sie Stream Data Archives. Für die Mechanik des Rollup-Recompute lesen Sie Rollups & recompute; das Streamdaten-Terminologie-Glossar definiert Tick, Bucket, Watermark, Generation und den Rest.

Ein Archiv lebt in zwei Stores​

Ein Archiv ist kein einzelnes Objekt in einer einzelnen Datenbank. Es ist über die beiden Stores aufgeteilt, die OctoMesh ohnehin verwendet, danach, worin jeder Store gut ist:

StoreWas er für ein Archiv enthält
MongoDB (Runtime-Graph)Die Definition des Archivs (Ziel-CK-type, Spaltenliste, Bucket-/Aggregationskonfiguration) und seine Runtime-State-Attribute (Status, den Rollup-Watermark, das Dirty-Window-Ledger, die Recompute-Health). Das Archiv ist eine Runtime-Entität – eine Instanz des System.StreamData-CK-Modells.
CrateDB (Zeitreihenebene)Die Massen-Zeitreihenzeilen – eine Pro-Tenant-, Pro-Archiv-Tabelle, die potenziell Milliarden von Datenpunkten hält.

Dies ist das Wichtigste, das man verinnerlichen muss: Die Runtime-/Mongo-Seite ist die Source of Truth für Identität und Zustand des Archivs, während die Zeitreihenzeilen in CrateDB leben. Der Lifecycle-Service hält die beiden ohne verteilte Transaktion abgleichbar (siehe Reconciliation in der Referenz).

Vergleichen Sie dies mit Runtime-Abfragen, die den aktuellen Zustand von Entitäten direkt aus MongoDB lesen. Streamdaten-Abfragen lesen Historie aus CrateDB, die über ihre rtId zurück auf die Entität verweist, die sie erzeugt hat.

Die drei Archivtypen​

Der abstrakte Basistyp System.StreamData/Archive hat drei konkrete Subtypen. Sie unterscheiden sich in ihrer Zeitachse und – entscheidend – darin, woher ihre Zeilen kommen:

Typ (CK type)ZeitachseZeilen erzeugt durch
System.StreamData/RawArchiveein timestamp pro Zeileexterne Producer (Pipelines, SDK) – Momentaufnahmen von Messungen
System.StreamData/TimeRangeArchiveexplizites [from, to) pro Zeileexterne Producer – bereits voraggregierte Daten (Smart-Meter-Intervalle, Wetter-APIs)
System.StreamData/RollupArchiveein Bucket [bucketStart, bucketEnd) pro Zeileausschließlich der System-Rollup-Orchestrator – wird niemals extern geschrieben

Ein Rollup aggregiert aus einem Quell-(Basis-)Archiv. Die Quelle kann ein Raw-Archiv, ein Time-Range-Archiv oder ein weiteres Rollup sein – Rollups bilden also Ketten (Rollup-of-Rollup). Ein Tages-Rollup könnte 15-Minuten-Rohdatenpunkte lesen; ein Jahres-Rollup könnte wiederum dieses Tages-Rollup lesen. Raw- und Time-Range-Archive akzeptieren externe Daten; Rollup-Spalten werden serverseitig erzeugt und direkte Schreibvorgänge werden abgelehnt.

Siehe Archive Subtypes für die Ingestion-Pfade und Anwendungsfälle je Typ.

Wie Zeilen gespeichert werden​

Jedes Archiv wird auf genau eine CrateDB-Tabelle im Schema des Tenants abgebildet. Die Spalten der Tabelle werden je nach Typ unterschiedlich abgeleitet:

  • Raw-/Time-Range-Archive – Spalten werden aus den konfigurierten CK-Attributpfaden abgeleitet, z. B. Amount.Value, ObisCode. Jeder erfasste Pfad wird zu einer echten, typisierten CrateDB-Spalte; verschachtelte Pfade und Arrays bilden auf OBJECT-/ARRAY-Spalten ab.
  • Rollup-Archive – Spalten werden aus den konfigurierten Aggregationen durch den RollupColumnGenerator abgeleitet, z. B. amountvalue_sum, dataquality_max. Es gibt genau eine Zeile pro Bucket pro Entität.

Jede Tabelle trägt außerdem die standardmäßigen Bookkeeping-Spalten: TargetCkTypeId (den konkreten CK type, relevant unter Vererbung), und Zeilen werden über ihre Quell-rtId zurück auf die erzeugende Entität verwiesen. Spalten können pro Spalte als indiziert oder nicht indiziert markiert werden. Rollup-Tabellen tragen zusätzlich eine generation-Spalte – Teil des Primärschlüssels –, die reader-sicheren Recompute erst möglich macht (siehe Lesepfad unten).

Das Schema ist nach der Aktivierung eingefroren

Sobald ein Archiv Activated ist, sind sein TargetCkTypeId und sein Spaltensatz unveränderlich – eine Schemaänderung bedeutet ein neues Archiv. Dies erlaubt es der Datenebene, für jeden Insert und jede Abfrage eine stabile Form anzunehmen. Siehe Storage Layout für die vollständigen Spalten-Mapping-Regeln.

Der Lebenszyklus​

Ein Archiv durchläuft die CkArchiveStatus-Zustände. Der Status ist es, der die gesamte Datenebene steuert – Inserts und Abfragen werden nur akzeptiert, solange ein Archiv Activated ist:

StatusCrateDB-TabelleLesen / Schreiben
Creatednoch nicht provisioniertabgelehnt – Definition existiert, keine Tabelle
Activatedprovisioniert, Schema eingefrorenakzeptiert
Disablederhalten (bewahrt)abgelehnt – erneut aktivieren, um fortzusetzen
Failedmöglicherweise partiellabgelehnt – Aktivierungs-DDL fehlgeschlagen; Wiederholung erforderlich

Aktivierung ist der Punkt, an dem das CrateDB-DDL läuft (CREATE TABLE) und das Schema einfriert. Sie füllt bewusst keine Historie nach: Der Watermark eines Rollups beginnt zum Aktivierungszeitpunkt, sodass nur Daten ab der Aktivierung automatisch gerollt werden. Das Befüllen der Historie ist eine explizite Betreiberaktion (ein Recompute über den historischen Bereich). Ein impliziter Scan der gesamten Historie bei der Aktivierung wäre eine unbegrenzte, überraschende Operation.

Siehe Lifecycle für das vollständige Zustandsdiagramm (einschließlich Lösch- und Wiederholungsübergänge) und das dreistufige Gate Activation Layers (Instanz → Tenant → Archiv).

Der Schreibpfad​

Externe Daten erreichen ein Raw- oder Time-Range-Archiv über die Mesh-Adapter-Pipeline-Nodes:

  • SaveStreamDataInArchive@1 → ein RawArchive
  • SaveTimeRangeStreamDataInArchive@1 → ein TimeRangeArchive

(Dieselbe Datenebene ist aus dem SDK über IStreamDataRepository erreichbar.) Bei jedem Schreibvorgang geschehen zwei Dinge, die für die Architektur bedeutsam sind:

  1. Ein Source-rtId-Integritäts-Guard stellt sicher, dass jede Zeile einer realen erzeugenden Entität zugeordnet ist, bevor sie persistiert wird.
  2. Der Hook zur Erkennung rückwirkender Schreibvorgänge läuft. Jeder Schreibvorgang wird gegen den Consumed Watermark etwaiger abhängiger Rollups geprüft: Ein Datenpunkt, der vor einem Fenster landet, das ein Rollup bereits aggregiert hat, wird als Dirty Window markiert, sodass er später neu berechnet werden kann. Dies ist rein Timestamp-vs-Watermark und wertunabhängig.

Dieser Erkennungsschritt ist der Einstiegspunkt in die gesamte Recompute-Maschinerie – siehe How a change in a base archive is detected, wie genau das Kriterium funktioniert.

Der Lesepfad​

Historie wird über Streamdaten-Abfragen gelesen, die in vier Formen verfügbar sind – simple, aggregation, grouping und downsampling – jeweils als persistierte (gespeicherte) oder transiente (Ad-hoc)-Variante. Der Überblick Stream Data Access behandelt die Abfrageoberfläche vollständig.

Für Rollup-Archive gibt es einen zusätzlichen Schritt, den der Leser nie sieht: Der Lesepfad injiziert den aktiven Generation-Filter in jede Abfrage, sodass ein Leser stets einen einzigen, konsistenten Snapshot jedes Fensters auflöst – entweder die alten Werte oder die neu berechneten, niemals einen halb geschriebenen Mix. Das macht Recompute reader-sicher: Eine Korrektur wird neben dem Live-Stand in eine neue Generation berechnet und mit einem einzigen Pointer-Flip committet. Siehe Optimistic atomic swap für den Mechanismus.

Zusammengesetzt​

Die Definition und der Zustand bleiben in Mongo; die Zeilen bleiben in CrateDB; Schreibvorgänge fließen in die Quelltabelle und kaskadieren in die Rollups; Lesevorgänge lösen eine konsistente Generation auf.

Wohin das führt​

Siehe auch​