Zum Hauptinhalt springen

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.

Archivdaten sichern

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 A definiertes 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 Activated ist.

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​

EigenschaftBeschreibung
Pro TenantJeder Tenant hat sein eigenes CrateDB-Schema. Keine tenantübergreifenden Daten.
TypisiertJeder erfasste Attributpfad wird zu einer echten CrateDB-Spalte mit deklariertem Typ.
Unveränderliches SchemaNach der Aktivierung können TargetCkTypeId und Columns nicht mehr geändert werden. Neue Version = neues Archiv.
Multi-ArchivEin CK-Typ kann von beliebig vielen Archiven gleichzeitig erfasst werden.
Status-gesteuertInserts und Abfragen werden nur akzeptiert, solange der Status des Archivs Activated ist.
Soft-Delete bewahrt DatenDas 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):

SubtypZeitachseIngestionspfadAnwendungsfall
RawArchiveEin timestamp pro ZeileExterne Producer (SaveStreamDataInArchive@1, IStreamDataRepository.InsertAsync)Momentanmessungen (Sensormesswerte, Maschinentelemetrie, Ereignisse)
RollupArchiveGebucketet [bucketStart, bucketEnd)Nur der System-Rollup-Orchestrator – keine externen InsertsAbgeleitete gebucketete Aggregationen eines anderen Archivs (raw oder Rollup-of-Rollup)
TimeRangeArchiveExpliziter [from, to)-Bereich pro ZeileExterne 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​

ElementKonvention
Schemaclean(tenantRtCkId.SemanticVersionedFullName), max. 63 Zeichen
Tabelleclean(targetCkType.SVF) + "_" + clean(archiveRtCkId.SVF), max. 200 Zeichen
SpaltecamelCase 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)​

SpalteTypHinweise
rtIdTEXTInstanz-ID der Entität
timestampTIMESTAMP WITH TIME ZONEZeitstempel des Datenpunkts (UTC)
ckTypeIdTEXTKonkreter Typ (relevant bei Vererbung)
rtCreationDateTimeTIMESTAMP WITH TIME ZONEDEFAULT CURRENT_TIMESTAMP
rtChangedDateTimeTIMESTAMP WITH TIME ZONEBei Konflikt aktualisiert
rtWellKnownNameTEXTOptional

Primärschlüssel: (timestamp, rtId, ckTypeId).

TimeRangeArchive ersetzt timestamp durch bucketStart + bucketEnd; der PK ist (bucketStart, bucketEnd, rtId, ckTypeId).

Pfad-zu-Spalte-Zuordnung​

PfadformCrateDB-SpaltentypNullbarkeit
Skalar (voltage)abgeleitet vom primitiven Typ des CK-AttributsNOT NULL wenn required
Record (sensor)OBJECT(STRICT) mit Unterfeldern aus CkRecordwie 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​

StatusTabellenzustandInserts / AbfragenHinweise
CreatedNicht provisioniertAbgelehntVollständig editierbar. Keine CrateDB-Ressourcen allokiert.
ActivatedProvisioniert, eingefrorenAkzeptiertSchema unveränderlich. Validierung auf dem Update-Pfad.
DisabledProvisioniert, eingefrorenAbgelehntDaten bewahrt. Zum Fortsetzen erneut aktivieren.
FailedKann partiell seinAbgelehntAktivierungs-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
EbeneQuelle
InstanzStreamData:Enabled in appsettings (Standard false)
TenantStreamDataGlobalSettings.IsEnabled (Flag pro Tenant)
ArchivCkArchive.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:

RolleErlaubte Operationen
StreamDataAdminVollständiger Lebenszyklus: Archive erstellen / aktualisieren / löschen; aktivieren / deaktivieren / einschalten / erneut versuchen; Tenant aktivieren / deaktivieren; Rollup-Lebenszyklus (Freeze / Unfreeze / Rewind).
StreamDataWriterIn aktivierte Archive einfügen. Metadaten lesen (Archive auflisten / abrufen). Kein Lebenszyklus, kein CRUD.
StreamDataReaderZeitreihenabfragen (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:

AlignmentBucket-Grenze
FixedSizeLastAggregatedBucketEnd + N · BucketSizeMs (offset-invariant; ignoriert die Zone)
CalendarDayLokale Mitternacht in ReferenceTimeZone
Iso8601WeekMontag 00:00:00 Lokalzeit in ReferenceTimeZone
CalendarMonthErster Tag des Monats 00:00:00 Lokalzeit in ReferenceTimeZone
CalendarYear1. 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.

Alignment und Zone sind nach der Aktivierung unveränderlich

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:

FunktionGespeicherte SpaltenFormel zur Lesezeit
Sum{base}_sumunverändert
Min{base}_minunverändert
Max{base}_maxunverändert
Count{base}_count (BIGINT)unverändert
Avg{base}_sum, {base}_count (beide gespeichert)sum / count – hält verkettete Rollups numerisch korrekt
First{base}_firstWert bei der frühesten Beobachtung im Bucket (numerische Quellspalten)
Last{base}_lastWert 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​

OperationMutationHinweise
CreatecreateRollupArchive(input)Der Server leitet TargetCkTypeId und Columns aus dem Quellarchiv + Aggregationen ab.
Bereich einfrierenfreezeRollupArchive(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.
UnfreezeunfreezeRollupArchive(rtId, acceptGaps)Idempotent. acceptGaps wird für die Auditierung erfasst, aber der Lückenerkennungs-Guard ist ein Follow-up.
Watermark zurückspulenrewindRollupWatermark(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 berechnenrecomputeArchive(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.

Deep Dive

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​

TriggerOberflächeVerhalten
ManuellGraphQL-Mutation recomputeArchive · POST /archives/{rollupRtId}/recompute · octo-cli RecomputeArchiveEin Operator erzwingt eine Recompute über [from, to), optional auf eine Entität eingegrenzt (rtIdScope). Erfordert StreamDataAdmin.
PeriodischRecomputeOrchestratorHostedService (Tick pro Tenant)Leert das Dirty-Windows-Ledger und berechnet nur die tatsächlich als dirty markierten Fenster neu – nicht das gesamte Archiv.
KettenpropagationinternNachdem 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:

  1. 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 immer generation = 0.
  2. 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.
  3. Eine winzige Side-Table pro Rollup (archive_<rtId>__genmap) hält die aktive Generation pro Bereich. Das Umschalten dieses Pointers auf N+1 ist ein Einzelzeilen-Schreibvorgang – der atomare Commit-Punkt.
  4. Der Lesepfad injiziert generation = CASE WHEN <range> THEN <activeGen> … ELSE 0 END, sodass jede Abfrage genau eine konsistente Generation pro Fenster auflöst.
  5. 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, Double oder DateTime (Ticks). String und nicht-skalare Typen werden nicht unterstützt. Ein Formel- / NaN- / NULL-Operand-Fehler speichert NULL für diese Zelle; die Ingestion schlägt nie fehl.
  • CK-Modell – CkArchiveColumn trägt Name + Formula + ResultType + den Engine-verwalteten ComputedState (seit System.StreamData 1.5.0) sowie ComputedVersion / PendingFormula (seit 1.6.1 bzw. 1.6.2). Eine Spalte mit einer Formula ist berechnet; eine mit einem Path wird 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 Pending persistiert, 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 auf Active. Ein Backfill-Fehler hinterlässt sie Failed und 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:

MutationLiefertBeschreibung
activateArchive(rtId)ArchiveTransitionResultProvisioniert die CrateDB-Tabelle, wechselt zu Activated.
disableArchive(rtId)ArchiveTransitionResultActivated → Disabled. Tabelle bleibt erhalten.
enableArchive(rtId)ArchiveTransitionResultDisabled → Activated. Validiert Spaltenpfade erneut.
retryArchiveActivation(rtId)ArchiveTransitionResultFailed → Activated. Idempotent (CREATE TABLE IF NOT EXISTS).
deleteArchive(rtId)BooleanVerwirft die CrateDB-Tabelle und soft-löscht die Entität. Abgelehnt, wenn aktive Rollups darauf verweisen.
createRollupArchive(input)OctoObjectIdServer-abgeleitete TargetCkTypeId + Columns.
createTimeRangeArchive(input)OctoObjectIdOperatordefinierter Zieltyp und operatordefinierte Spalten.
freezeRollupArchive(rtId, until)ArchiveTransitionResultMonotones FrozenUntil setzen.
unfreezeRollupArchive(rtId, acceptGaps)ArchiveTransitionResultSetzt FrozenUntil zurück.
rewindRollupWatermark(rtId, toBucketEnd)ArchiveTransitionResultAggregiert ab der angegebenen Grenze neu.
recomputeArchive(rtId, from, to, rtIdScope?)RecomputeJobInfoLöst eine optimistische Recompute von [from, to) aus / verschmilzt sie, optional pro Entität. Liefert den Job-Snapshot.
addComputedColumn(rtId, name, formula, resultType, indexed?)ArchiveTransitionResultFügt eine Computed Column hinzu und füllt sie per Backfill (optimistisch / atomar).
updateComputedColumnFormula(rtId, name, formula)ArchiveTransitionResultFormuliert eine Computed Column neu (versionierter Backfill + atomarer Swap). Abgelehnt, wenn eine andere Computed Column sie referenziert.
removeComputedColumn(rtId, name)ArchiveTransitionResultEntfernt 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​

MethodePfadBeschreibung
POST/enable / /disableUmschalter auf Tenant-Ebene. /disable antwortet mit 409, solange noch irgendein Archiv Activated ist (AB#4255).
POST/archives/{archiveRtId}/activateWie GraphQL activateArchive.
POST/archives/{archiveRtId}/disable / /enable / /retryStatusübergänge.
DELETE/archives/{archiveRtId}Wie GraphQL deleteArchive.
POST/archives/{rollupRtId}/freeze / /unfreeze / /rewindNur-Rollup-Steuerungen.
POST/archives/{rollupRtId}/recomputeOptimistische Recompute von [from, to) (rtIdScope optional). Liefert den Job-Snapshot.
GET/archives/{archiveRtId}/recompute-jobsAktuelle Recompute-Jobs (neueste zuerst) zum Debuggen.
POST/archives/{archiveRtId}/insertTimeRangeBulk-Time-Range-Insert.
GET/archives/{archiveRtId}/rollupsAktive 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​

  1. Archiv auflösen (einzelner Lookup), Statusprüfung.
  2. Jeden Punkt im Batch vorvalidieren – Abdeckung erforderlicher Pfade, Typkompatibilität, Pfadauflösung. Stoppt beim ersten Verstoß.
  3. Bei Verstoß → mit dem betreffenden Index werfen. Keine Zeile geschrieben.
  4. 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.
  • RawRetentionMs auf dem Basistyp Archive definiert den Aufbewahrungshorizont. null = für immer aufbewahren.
  • Der ArchiveRetentionScheduler (IHostedService, stündlich) verwirft pro Archiv und pro Tenant Partitionen, die älter als now() - RawRetentionMs sind. 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:

  1. Alles DDL ist idempotent (IF EXISTS / IF NOT EXISTS).
  2. Die Operationsreihenfolge ist Crate-zuerst, Mongo-zuletzt: Activate führt CREATE TABLE aus, dann schaltet es Status in Mongo um. Delete führt DROP TABLE aus, dann soft-löscht es die Entität. Mongo ist die Quelle der Wahrheit.
  3. 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 mit ArchiveNotActivatedException ab.
  4. Beim Start zählt der ArchiveReconciler die Activated-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​

SzenarioVerhalten
Zwei parallele Activate auf demselben ArchivSerialisiert über RepositoryDistributedLockService. Der zweite Aufruf beobachtet Activated und kehrt idempotent zurück.
Disable, während Inserts laufenLaufende Inserts werden abgeschlossen. Neue Inserts nach dem Umschalten lösen ArchiveNotActivatedException aus.
Delete, während Inserts laufenWie Disable, plus das DROP TABLE kann einmal mit kurzem Backoff wiederholen, falls die Tabelle kurzzeitig gesperrt ist.
Activate schlägt mitten im DDL fehlStatus → Failed. EnsureArchiveCreatedAsync ist idempotent (CREATE TABLE IF NOT EXISTS), sodass ein Retry sicher ist.
Bulk-Insert auf deaktiviertem ArchivGesamter 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.

ExceptionGeworfen, wenn
StreamDataExceptionBasis – generischer StreamData-Fehler.
StreamDataNotEnabledExceptionInstanz- oder Tenant-Flag ist false bei einer Operation, die es erfordert.
ArchiveNotFoundExceptionarchiveRtId löst sich nicht zu einer Archive-Entität auf.
ArchiveNotActivatedExceptionArchiv existiert, aber Status ≠ Activated beim Insert / bei der Abfrage.
ArchiveSchemaImmutableExceptionEin Update nach der Aktivierung versucht, schemarelevante Felder zu ändern.
ArchivePathInvalidExceptionEin CkArchiveColumn.Path kann nicht gegen TargetCkTypeId aufgelöst werden.
ArchiveColumnTypeUnsupportedExceptionEin Attributtyp kann nicht auf einen CrateDB-Spaltentyp abgebildet werden.
RequiredAttributeMissingExceptionEin Insert hat keinen Wert für einen required: true-Skalarpfad. Trägt den Index des betreffenden Punkts.
ArchiveActivationFailedExceptionDie DDL-Ausführung während der Aktivierung ist fehlgeschlagen. Umschließt den zugrunde liegenden SQL-Fehler.
InvalidArchiveStateTransitionExceptionEine Lebenszyklus-Mutation wurde gegen einen inkompatiblen Status aufgerufen.
RollupSourceInUseExceptiondeleteArchive auf einem Raw-Archiv aufgerufen, das noch aktive Rollups hat.
RollupSourceMissingExceptioncreateRollupArchive verweist auf eine Quell-archiveRtId, die nicht existiert oder noch nicht aktiviert ist.

Observability​

Metriken​

MetrikTypTags
streamdata.archive.countGaugetenant, status
streamdata.archive.status_transitionsCountertenant, archive, from, to
streamdata.archive.activation_duration_msHistogramtenant, archive, outcome
streamdata.insert.duration_msHistogramtenant, archive, batch_size_bucket
streamdata.insert.pointsCountertenant, archive
streamdata.insert.required_violationsCountertenant, archive
streamdata.query.duration_msHistogramtenant, archive, query_type
streamdata.query.rows_returnedHistogramtenant, archive, query_type
streamdata.crate.connections.openGauge(instanzweit)
streamdata.crate.retention_partitions_droppedCountertenant, archive
streamdata.crate.rollup_rows_insertedCountertenant, 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​

  1. Ein Archiv pro Belang. Unterschiedliche Aufbewahrung, unterschiedliche Abfragemuster, unterschiedliche Konsumenten → unterschiedliche Archive.
  2. Seien Sie explizit bei required. Erforderliche Skalare werden zu NOT NULL-Spalten und werden app-seitig vor dem Insert validiert.
  3. Indizieren Sie sparsam. Indizes kosten Speicher und Insert-Durchsatz. Standardmäßig indiziert; verzichten Sie darauf für Spalten, die gelesen, aber niemals gefiltert werden.
  4. Rollen Sie früh auf. Ein RollupArchive ist über lange Horizonte weitaus günstiger abzufragen als ein Raw-Scan.
  5. Wählen Sie das Alignment für den Anwendungsfall. Energie- / EDA-Workflows wollen Kalender-Alignments (monatlich, wöchentlich); Engineering-Telemetrie will üblicherweise FixedSize.
  6. 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: false auf (oder spulen Sie die Watermark zurück), um sauber neu zu aggregieren.
  7. Zuerst soft-löschen. Deaktivieren Sie ein Archiv, um zu bestätigen, dass kein Konsument bricht, bevor Sie löschen.

Siehe auch​