Zum Hauptinhalt springen

Rollups & recompute

Diese Seite ist der an Ingenieure gerichtete Begleiter zu Stream Data Archives. Sie erklärt, wie ein RollupArchive seine abgeleiteten Aggregationen korrekt hält, während sich die zugrunde liegenden Daten ändern – das Erkennungskriterium, das Dirty-Windows-Ledger, der Abhängigkeitsgraph, der periodische Orchestrator und der optimistische atomare Swap, der Leser durchgehend auf einem konsistenten Snapshot hält.

Wenn Sie nur die an Betreiber gerichtete Anleitung benötigen (die Studio-Schaltflächen und wann man sie verwendet), lesen Sie stattdessen den Abschnitt Studio: Archives — Rollup-Specific Actions. Diese Seite befasst sich mit der Mechanik hinter jenen Schaltflächen.

Wo dies liegt
  • CK-Modell – System.StreamData (seit 1.6.2); die Rollup-Runtime-State-Attribute und die RecomputeJob-Entität sind seit dieser Modellversion verfügbar.
  • Engine – octo-construction-kit-engine + octo-construction-kit-engine-mongodb (Orchestratoren, Dirty-Detection, Abhängigkeitsgraph).
  • API & Persistenz – octo-asset-repo-services (GraphQL-/REST-Oberfläche, CrateDB-Generation-Tabellen, MongoDB-Runtime-State).
Terminologie

Das über diese Seite hinweg verwendete Vokabular – Tick, Bucket (sowie Bucket-Größe / -Ausrichtung), Watermark (sowie Watermark-Lag / Consumed Watermark), Dirty Window, Generation, Source vs Rollup Archive, Forward Aggregation / Recompute / Rewind, Chain Propagation und Coalesce – ist einmal im zentralen Streamdaten-Glossar definiert. Erstverwendungen unten verlinken direkt auf den relevanten Eintrag. Das durchgehende Beispiel ist der voest energy-Fall: 15-Minuten-Quelldatenpunkte, die zu Tages- / Stunden- / usw. Aggregationen gerollt werden.

Die drei Wege, auf denen Rollup-Daten (neu) berechnet werden​

Die gespeicherten Zeilen eines Rollups können durch genau drei Codepfade erzeugt oder ersetzt werden. Alles andere (das Dirty-Ledger, Chain Propagation, die Studio-Schaltflächen) treibt letztlich einen dieser drei an.

  1. Forward Aggregation. Der RollupOrchestrator tickt und aggregiert den nächsten geschlossenen Bucket [watermark, now - watermarkLag), einen Bucket pro Tick. Er schreibt einen Upsert bei generation 0 und rückt den Watermark vor (LastAggregatedBucketEnd). Dies ist der Steady-State-Pfad – er bewegt sich immer nur vorwärts.
  2. Recompute (manuell oder periodisch). Eine bereichsbeschränkte [from, to)-Neuaggregation, die als optimistischer atomarer Swap läuft: Sie ist reader-sicher (Abfragen liefern durchgehend einen konsistenten Snapshot) und erzeugt einen RecomputeJob-Datensatz. Dies ist der Pfad, um Historie ohne Ausfallzeit zu korrigieren.
  3. Rewind Watermark. Bewegt die Forward-High-Water-Mark rückwärts. Der Orchestrator aggregiert dann über den zurückgesetzten Bereich vorwärts neu und überschreibt den generation 0-Upsert an Ort und Stelle. Dies ist nicht reader-sicher – Leser können teilweise neu aggregierte Buckets beobachten – und daher destruktiv.

Zusätzlich entscheidet ein Mechanismus der automatischen Dirty-Detection, wann ein Recompute nötig ist. Er schreibt selbst keine Zeilen um; er zeichnet Intent auf (ein Dirty Window), das der periodische Orchestrator später in einen Recompute verwandelt.

Open-Bucket-Refresh (AB#4306)

Forward Aggregation schließt einen Bucket erst, sobald er vollständig in der Vergangenheit liegt (< now - watermarkLag), sodass eine grobe Teilperiode – dieser Monat bisher, dieses Jahr bisher – andernfalls keinen Wert zeigen würde, bis die Periode endet. Um jene Teilperioden-Summen live zu halten, aggregiert der RollupOrchestrator zusätzlich bei jedem Tick den aktuellen offenen (laufenden) Bucket neu und schreibt ihn bei generation 0, ohne den Watermark vorzurücken. Dies wird über die Option RefreshOpenBucket gesteuert (StreamData:Rollup:RefreshOpenBucket, Standard true) und überspringt einen Bucket, der eingefroren ist (siehe Freeze).

Weil der offene Bucket die exklusive Domäne des Forward-Pfads ist, schließt Recompute ihn bewusst aus: Ein manueller/periodischer Recompute begrenzt sein to auf den Beginn des Buckets, der now enthält, sodass nur vollständig geschlossene Buckets im Bereich verbleiben (ein vollständig offener Bereich wird zu einem Completed-No-op). Würde ein Recompute den offenen Bucket beanspruchen, würde er die genmap-Generation für dieses Fenster auf N+1 kippen, und der Lesepfad (höchste Generation gewinnt) würde dann den generation 0-Schreibvorgang des Refresh maskieren – wodurch die Teilperioden-Summe beim Recompute einfrieren würde, statt neue Quelldaten zu verfolgen. Das Begrenzen des Bereichs hält den offenen Bucket auf der stets frischen generation 0-Spur. Dies ist am wichtigsten in einer Cascade (Rollup-of-Rollup): Der offene Bucket jeder groben Ebene (täglich → monatlich → jährlich) wird günstig aus einer Handvoll Zeilen der darunterliegenden Ebene aufgefrischt, sodass die aktuelle Periode die ganze Leiter hinauf aktuell bleibt.

Vergleich​

AspektForward AggregationAutomatische Dirty-DetectionRecomputeRewind Watermark
Was es auslöstRollupOrchestrator-Tick (Steady State)Jeder Archiv-Schreibvorgang, über RetroactiveWriteDetectorrecomputeArchive (Manual) · RecomputeOrchestratorHostedService (Periodic) · Chain (ChainPropagation)rewindRollupWatermark (Betreiber)
Was es tutAggregiert den nächsten geschlossenen BucketZeichnet ein CkArchiveDirtyWindow auf; schreibt nichts umAggregiert [from, to) in eine neue Generation neu und kippt dann den PointerBewegt LastAggregatedBucketEnd zurück; der Forward-Pfad läuft erneut
Reader-sicherJa (Append/Upsert bei gen 0)Nicht zutreffend (keine Zeilenänderungen)Ja – konsistenter Alt-oder-Neu-Snapshot, niemals gemischtNein – Leser können partielle Neuaggregation sehen
ScopeEin Bucket pro Tick, nur vorwärtsEin [min ts .. max ts + 1 tick)-Fenster, begrenzt durch die Bounded Retro Reach-ObergrenzeEin begrenzter [from, to)-Bereich, optional per rtIdScope auf eine Entität beschränktAlles ab der Rewind-Grenze vorwärts
ObservabilityLastAggregatedBucketEndArchive.DirtyWindows / PendingRecomputeRangesRecomputeJob-Historie + Archive.LastRecompute*-FelderÄnderung von LastAggregatedBucketEnd
DestruktivNeinNeinNeinJa
Wann verwendenImmer aktiv; nichts zu tunImmer aktiv; nichts zu tunHistorische Daten korrigieren, während Dashboards/Pipelines live bleibenLetztes Mittel / Bulk-Backfills, bei denen Leser eine Lücke tolerieren können
Faustregel

Bevorzugen Sie Recompute für Korrekturen. Greifen Sie zu Rewind nur, wenn Sie bewusst einen großen Forward-Bereich erneut ausführen (zum Beispiel nach einem Bulk-Backfill in das Quellarchiv) und tolerieren können, dass Leser kurzzeitig in Bearbeitung befindliche Buckets sehen.

Wie eine Änderung in einem Basisarchiv erkannt wird​

Dies ist die Frage, die Ingenieure am häufigsten stellen: Wenn niemand dem Rollup mitteilt „Zeile X hat sich geändert", woher weiß es, dass es neu berechnen muss? Die Antwort lautet: rein Timestamp-vs-Watermark zur Insert-Zeit, und sie ist wertunabhängig.

Bei jedem Schreibvorgang in ein Archiv läuft RetroactiveWriteDetector.TryBuildDirtyWindow. Für das geschriebene Archiv berechnet es den Consumed Watermark:

consumedWatermark = das Maximum von LastAggregatedBucketEnd über die nicht gelöschten abhängigen Rollups des Quellarchivs – d. h. der am weitesten fortgeschrittene Punkt, den irgendein Abhängiger bereits hinter sich aggregiert hat.

Warum das Maximum, nicht das Minimum (AB#4288)

Eine Quelle kann Rollups sehr unterschiedlicher Granularitäten speisen, deren Watermarks weit auseinanderliegen (ein Jahres-Rollup hinkt hinterher – sein Bucket schließt nur jährlich – während ein Stunden-Rollup nahe now liegt). Die Verwendung des Minimums würde eine Korrektur stillschweigend übersehen, die im Band zwischen dem langsamsten und dem schnellsten Abhängigen landet. Das Maximum zu nehmen ist eine bewusste sichere Überannäherung: Es markiert einen solchen Schreibvorgang als rückwirkend; der Orchestrator klemmt dann den Recompute-Bereich jedes Abhängigen auf dessen eigenen Watermark (AB#4288), sodass ein hinterherhinkender Abhängiger niemals einen noch nicht geschlossenen Bucket neu aggregiert.

Dann, für jeden Datenpunkt-Timestamp ts im Schreibvorgang:

  • Wenn consumedWatermark null ist (noch kein Abhängiger hat etwas konsumiert) → nichts ist dirty. Es gibt keine Aggregation zu invalidieren.
  • Wenn ts >= consumedWatermark → der Schreibvorgang ist ein Forward-Append und wird vom Detector ignoriert. Der normale Forward-Aggregation-Pfad wird ihn aufgreifen.
  • Wenn ts < consumedWatermark → der Schreibvorgang ist ein RetroactiveModify: Er landet in einem Fenster, das ein Abhängiger bereits aggregiert hat, sodass diese Aggregation nun veraltet ist. Der Detector zeichnet ein CkArchiveDirtyWindow auf, das [min(retroactive ts) .. max(retroactive ts) + 1 tick) umspannt.
Grenzen-Nuance (< nicht <=)

Die Erkennung ist strikt ts < consumedWatermark. Ein Schreibvorgang, dessen Timestamp genau auf dem Watermark landet, wird als forward (Append) behandelt, nicht als rückwirkend. Behalten Sie dies im Hinterkopf, wenn Sie über Grenzfälle an einer Bucket-Grenze nachdenken.

Es vergleicht keine Werte​

Der Detector vergleicht niemals Werte. Eine idempotente Neu-Ingestion byte-für-byte identischer Daten in ein bereits konsumiertes Fenster markiert dieses Fenster dennoch als dirty. Dies ist eine bewusste sichere Überannäherung: Sie kann eine echte Änderung niemals verpassen, um den Preis, gelegentlich einen Recompute einzuplanen, der identische Ausgabe erzeugt.

Wertunabhängig by Design

Eine zukünftige No-op-Optimierung (den Recompute überspringen, wenn die neu ingestierten Werte unverändert sind oder unter einem Delta-Schwellenwert liegen) ist denkbar, aber noch nicht als eigenes Item erfasst. Bis dahin behandeln Sie jeden rückwirkenden Schreibvorgang als garantiertes Dirty Window.

Threshold-Reset-Policy – AB#4196

Ob eine erkannte rückwirkende Änderung den „computed up to"-Schwellenwert manuell oder automatisch zurücksetzen soll, ist Gegenstand von AB#4196. Die Antwort lautet beides, und der automatische Pfad wird bereits ausgeliefert: Die automatische Erkennung zeichnet ein Dirty Window auf und treibt einen bereichsbeschränkten Recompute (sie läuft nicht den Forward-Watermark zurück), während rewindRollupWatermark als manuelle Bulk-Notluke bestehen bleibt. Die automatische Reichweite ist durch eine Pro-Archiv-Obergrenze begrenzt – siehe Bounded retro reach unten. Engine-Entscheidungsdokument: octo-construction-kit-engine/docs/concept-rollup-threshold-reset.md.

Bounded retro reach (AB#4196)​

Der automatische Pfad darf niemals zulassen, dass ein einzelner sehr später Schreibvorgang einen unbegrenzten Recompute nach sich zieht. Stellen Sie sich eine Korrektur vor, die mit einem fünf Jahre alten Timestamp eintrifft: Ohne Guard würde sie ein Fünf-Jahres-Fenster als dirty markieren und einen Recompute jedes Buckets bis zu diesem Punkt zurück einplanen. Die Bounded-Retro-Reach-Obergrenze verhindert das.

  • Config. Archive.MaxRetroactiveReachMs (pro Quellarchiv, seit System.StreamData 1.6.8; null = unbegrenzt, das Verhalten vor 1.6.8) plus eine Flotten-Obergrenze StreamData:Recompute:MaxRetroactiveReachHardLimitMs in der Host-Konfiguration. Die effektive Obergrenze, die bei der Erkennung angewendet wird, ist min(pro-Archiv, Flotten-Obergrenze) – ein Pro-Archiv-Wert kann die Obergrenze nur verschärfen, niemals lockern.
  • Wo sie greift. Primär bei der Erkennung: Der RetroactiveWriteDetector begrenzt das automatische Dirty Window nach unten auf consumedWatermark - effectiveCap. Ein rückwirkender Timestamp älter als dieser Boden wird aus dem automatischen Fenster verworfen; ist der gesamte rückwirkende Batch älter als der Boden, wird kein Dirty Window aufgezeichnet. In jedem Fall wird ein WARN geloggt (retroactive write reached beyond the automatic recompute cap …), sodass der Verwurf sichtbar ist. Als doppelte Absicherung wendet der periodische Orchestrator die Pro-Archiv-Obergrenze der Quelle erneut an, wenn er ein Dirty Window auf jeden Abhängigen propagiert (er begrenzt den Start des Stale-Bereichs nach unten auf dependentWatermark - cap), was auch jedes vor Existenz der Obergrenze aufgezeichnete Dirty Window begrenzt.
  • Was sie nicht begrenzt. Die Obergrenze ist nur automatisch. recomputeArchive und rewindRollupWatermark bleiben unbegrenzt – sie sind die bewusste Betreiber-Notluke für eine echte tiefe Korrektur, die außerhalb der automatischen Reichweite liegt.
Standard

Wird pro Archiv als null (unbegrenzt) ausgeliefert und bewahrt so das Verhalten vor 1.6.8 – kein bestehender Tenant verliert stillschweigend eine legitime Monate alte Korrektur. Setzen Sie sie pro Archiv (z. B. P30D–P90D für eine hochfrequente Energiequelle) und/oder setzen Sie die Flotten-Obergrenze in der Produktions-Host-Konfiguration konservativ, um einen pathologischen Runaway zu begrenzen.

Das Dirty-Dependents-Ledger und Chain Propagation​

Die Erkennung zeichnet Intent auf; Propagation und Recompute handeln danach.

  • Ledger-Datensätze. CkArchiveDirtyWindow (ein erkanntes veraltetes Fenster auf einem Quellarchiv) und CkArchiveRecomputeRange (ein für einen bestimmten Abhängigen eingereihter Bereich) sind die beiden Ledger-Datensatztypen. Sie leben als Runtime-State-Attribute auf der Archiv-Entität (siehe State & observability), nicht als CrateDB-Spalten.
  • Abhängigkeitsgraph. RollupDependencyGraph durchläuft die transitive Menge der Abhängigen, sodass eine mehrstufige Rollup-of-Rollup-Kette top-down behandelt wird. Der Durchlauf ist zyklensicher – eine fehlkonfigurierte zyklische Abhängigkeit kann die Propagation nicht in eine Endlosschleife schicken.
  • Periodischer Orchestrator. RecomputeOrchestratorHostedService ist ein BackgroundService (Standard-Tick 60 s, in DI registriert). Bei jedem Tick propagiert er ausstehende Dirty Windows in die PendingRecomputeRanges jedes betroffenen Abhängigen und berechnet dann nur jene Bereiche neu – niemals das gesamte Archiv. Diese Läufe tragen den Trigger Periodic.
  • Chain Propagation. Nachdem der Recompute eines Rollups erfolgreich committet wurde, werden die direkten Abhängigen dieses Rollups erneut eingereiht, diesmal mit dem Trigger ChainPropagation. Dies ist es, was eine Korrektur durch eine Rollup-of-Rollup-Kette vorwärtsrollt.
  • Coalescing. Überlappende Trigger für dasselbe Archiv koaleszieren – ihre Bereiche verschmelzen zu einem einzigen aktiven Job und der verdrängte Trigger wird als Coalesced aufgezeichnet. Es gibt höchstens einen aktiven Recompute pro Archiv zur selben Zeit.

End-to-End-Fluss​

Der vollständige Pfad von einem späten Schreibvorgang zu einer vollständig propagierten Korrektur:

Leser, die das Rollup zu einem beliebigen Zeitpunkt abfragen, lösen ihre Generation über die genmap auf (siehe unten) und sehen daher entweder den alten Snapshot oder den neuen – niemals einen halb geschriebenen Mix.

Optimistischer atomarer Swap​

Der Recompute wird intern als „stabil, optimistisch" beschrieben: stabil, weil Leser niemals blockiert werden oder ein partielles Ergebnis sehen, optimistisch, weil er die neue Antwort neben der alten berechnet und erst mit einem einzigen Pointer-Flip committet.

Die Maschinerie:

  • Jede Rollup-CrateDB-Tabelle trägt eine generation BIGINT-Spalte, die Teil des Primärschlüssels ist. Forward Aggregation schreibt immer generation 0.

  • Eine Nebentabelle, archive_<rtId>__genmap, hält die aktive Generation pro Fensterbereich.

  • Ein Recompute berechnet seine Ausgabe in eine Pro-Job-Staging-Tabelle bei Generation N+1, kopiert diese Zeilen in die Live-Tabelle (sodass beide Generationen kurz koexistieren) und kippt dann den genmap-Pointer auf N+1. Dieser Ein-Zeilen-Flip ist der atomare Commit-Punkt. Schließlich fegt er die verdrängten Generationen weg.

  • Der Lesepfad injiziert die aktive Generation in jede Abfrage:

    generation = CASE WHEN <range> THEN <gen> ELSE 0 END

    sodass ein Leser transparent Generation 0 außerhalb des neu berechneten Bereichs und die aktive neu berechnete Generation innerhalb dessen aufgreift.

Konsequenzen:

  • Leser sehen stets einen konsistenten Alt-oder-Neu-Snapshot, niemals gemischt.
  • Ein Crash vor dem Flip lässt die vorherigen Werte vollständig intakt – die Staging-Zeilen und nicht referenzierten gen N+1-Zeilen sind nur Müll, der weggefegt wird; nichts, was der Leser sehen kann, hat sich geändert.
  • Genau deshalb ist Recompute reader-sicher und Rewind nicht: Rewind überschreibt generation 0 an Ort und Stelle, ohne eine zweite Generation, auf die zurückgefallen werden kann.

Referenz-Enums​

EnumWerteBedeutung
CkRecomputeTriggerManual · Periodic · ChainPropagationWarum ein Recompute-Job existiert: betreiberangefordert, vom periodischen Orchestrator aus dem Dirty-Ledger abgearbeitet oder von einem erfolgreichen vorgelagerten Recompute eingereiht.
CkRecomputeJobStatePending · Running · Swapping · Completed · Failed · CoalescedJob-Lebenszyklus. Swapping ist das genmap-Flip-Fenster; Coalesced bedeutet, dass der Job in einen überlappenden aktiven Job verschmolzen wurde.
Change kindAppend · RetroactiveModifyKlassifikation eines einzelnen Schreibvorgangs durch RetroactiveWriteDetector: forward (ignoriert) vs. rückwirkend (zeichnet ein Dirty Window auf).
Change sourceManual · Pipeline · ImportHerkunft des Schreibvorgangs, der die Erkennung ausgelöst hat – für Audits offengelegt.

State & Observability​

Der gesamte Rollup-Recompute-Zustand liegt als Runtime-State-Attribute auf der Archiv-Entität in MongoDB sowie in der RecomputeJob-Entitätshistorie – nicht als CrateDB-Spalten (die Generation-Spalte und genmap sind ein Implementierungsdetail der Datenebene, nicht der beobachtbare Zustand).

Attribut auf ArchiveWas es enthält
Archive.DirtyWindowsErkannte veraltete Fenster, die noch nicht neu berechnet wurden
Archive.PendingRecomputeRangesFür den nächsten Recompute dieses Archivs eingereihte Bereiche
Archive.RecomputeInProgressOb aktuell ein Recompute läuft
Archive.LastRecomputeStartedAtBeginn des jüngsten Recompute
Archive.LastRecomputeSuccessAtAbschluss des jüngsten erfolgreichen Recompute
Archive.LastRecomputeFailureAtZeitpunkt des jüngsten Fehlschlags
Archive.LastRecomputeFailureReasonGrund-String für den jüngsten Fehlschlag

Diese (und die Pro-Job-Historie) werden offengelegt über:

  • GraphQL – rollupsFor(archiveRtId) für die Live-Pro-Rollup-Health-Felder und recomputeJobsFor(archiveRtId) für die Job-Historie (neueste zuerst).
  • REST – GET /archives/{archiveRtId}/rollups und GET /archives/{archiveRtId}/recompute-jobs.
  • Refinery Studio – das Rollups-Panel in der Archiv-Detailansicht (die Recompute-Spalte und die Recompute-Job-Historientabelle).

Trigger und Befehle​

OperationGraphQLocto-cli
Einen Bereich neu berechnenrecomputeArchive(rtId, from, to, rtIdScope?)RecomputeArchive
Recompute-Jobs auflistenrecomputeJobsFor(archiveRtId)ListRecomputeJobs
Freeze / UnfreezefreezeRollupArchive / unfreezeRollupArchiveFreezeRollupArchive / UnfreezeRollupArchive
Watermark zurücksetzen (Rewind)rewindRollupWatermark(rtId, toBucketEnd)RewindRollupWatermark
Aktivierung füllt keine Historie nach

Das Aktivieren eines Rollups aggregiert nicht rückwirkend die bestehende Historie des Quellarchivs – der Watermark beginnt bei now - bucketSize, sodass nur Daten ab der Aktivierung automatisch gerollt werden. Das Befüllen der Historie ist eine explizite Betreiberaktion (ein recomputeArchive über den historischen Bereich oder ein Rewind). Dies ist beabsichtigt: Ein impliziter Scan der gesamten Historie bei der Aktivierung wäre eine unbegrenzte, überraschende Operation.

Lebenszyklus-Interaktionen: Aktivierung, Deaktivierung, Löschung (AB#4300)​

Eingereihte Recompute-Arbeit wird immer nur für ein Activated-Rollup abgearbeitet – der periodische Orchestrator überspringt jedes Rollup, das nicht Activated ist. Zwei Guards verhindern, dass daraus ein Job entsteht, der niemals Fortschritt machen kann:

  • Fail-fast auf einem nicht aktivierten Rollup. Ein backfillRollup oder recomputeArchive gegen ein Rollup, das Created, Disabled oder Failed ist, gibt sofort einen Failed-Job zurück, mit einem Grund wie „Archive is not activated (status Disabled) — activate it before backfilling." – statt einen Pending-Job + Bereich vorab anzulegen, den die Abarbeitung stillschweigend für immer überspringen würde. Aktivieren Sie das Rollup zuerst, reihen Sie dann die Arbeit ein.
  • Purge bei Deaktivierung / Löschung. Das Deaktivieren oder Löschen eines Archivs purged seine eingereihte Recompute-Arbeit: Die PendingRecomputeRanges werden geleert und der einzelne aktive (Pending/Running/Swapping)-Job wird mit einem erklärenden Grund als Failed beendet. Dies verhindert einen verbleibenden, nicht verarbeitbaren Pending-Job und stoppt, dass ein Delete-dann-Re-Import, der dieselbe Runtime-ID wiederverwendet, veraltete Bereiche/Jobs erbt.
Bestehende Geister heilen sich selbst

Ein Pending-Job, der vor dem Deaktivieren eines Rollups eingereiht wurde (d. h. unter einer früheren Engine-Version, ohne diese Guards), wird nicht rückwirkend umgeschrieben – aber er heilt sich selbst, sobald das Rollup erneut aktiviert wird: Der nächste Orchestrator-Tick adoptiert jenen Pending-Job und treibt ihn zu Completed.

Siehe auch​