Dauerhafte Funktionen mit pg_durable für Azure HorizonDB (Vorschau)

pg_durable ist das dauerhafte Ausführungsmodul in Azure HorizonDB. Sie können lange ausgeführte, mehrstufige SQL-Workflows (Einbettungspipelinen, ETL-Aufträge, KI-Aufrufe, geplante Aufträge, Genehmigungsflüsse) definieren und mit den gleichen Zuverlässigkeitsgarantien ausführen, die Sie von einem dedizierten Orchestrator wie Durable Functions erwarten, ohne Postgres verlassen zu müssen.

pg_durable ist auch die Ausführungsschicht, die robusten KI-Pipelines zugrunde liegt. Wenn Sie KI-Pipelines verwenden, sorgt pg_durable dafür, dass sie Abstürze überstehen, bei Fehlern erneut versuchen und ab dem letzten abgeschlossenen Schritt fortgesetzt werden können.

Note

pg_durable befindet sich in der Vorschau.

Was "dauerhaft" bedeutet

Eine permanente Funktion in pg_durable wird bei jedem Schritt auf dem Datenträger gespeichert. Dadurch erhalten Sie eine bestimmte Reihe von Garantien, die Sie nicht von einem einfachen BEGIN ... COMMIT Block oder einem Cron-Job erhalten:

  • Überdauert Datenbankabstürzen und Neustarts. Abgeschlossene Schritte werden nicht erneut ausgeführt, wenn der Server wieder hochfährt. Bereits begonnene Schritte werden ab dem letzten Checkpoint fortgesetzt. Ausstehende Schritte werden ausgeführt, wenn der Worker wieder online ist.
  • Überlebt lange Wartezeiten. Ein Workflow kann stundenlang schlafen, auf einen cron-Zeitplan warten oder ein externes Signal blockieren und an der Stelle weitergehen, an der er unterbrochen wurde.
  • Übersteht Ausfälle. Fehlgeschlagene Schritte können automatisch wiederholt werden, ohne die gesamte Funktion erneut auszuführen.
  • Erfasst Identitätsdaten. Eine Funktion wird mit den Berechtigungen des Benutzers ausgeführt, der ihn gestartet hat, nicht mit den Berechtigungen des Hintergrundarbeiters. Mehrinstanzenfähige Workloads bleiben isoliert.
  • Bleibt in SQL sichtbar. Sie können Status, Verlauf, Ausführungsanzahl und Ausgaben über dieselbe Schnittstelle prüfen, die Sie für alles andere in HorizonDB verwenden: eine SELECT Anweisung.

Was die Dauerhaftigkeit nicht automatisch tut: Sie macht keine nichtidempotenten externen Vorgänge sicher, um den Vorgang selbst zu wiederholen. Wenn ein Schritt eine externe API aufruft, die Geld belastet, entwerfen Sie den Schritt als idempotent (z. B. durch Übergeben eines idempotenten Schlüssels).

Wann pg_durable verwendet werden sollte

Verwenden Sie pg_durable, wenn Sie damit arbeiten müssen:

  • Es dauert lange genug, bis es in der Mitte zu einem Fehler kommt (Einbettung von Daten über Millionen von Zeilen, ein mehrstufiger ETL-Auftrag, ein Backfill).
  • Im Fehlerfall muss der Vorgang erneut versucht werden, ohne die bereits erfolgreich abgeschlossenen Teile erneut auszuführen.
  • Muss nach einem Zeitplan ausgeführt werden (jede Stunde, werktags um 9 Uhr).
  • Muss auf ein externes Ereignis warten (eine Genehmigung, ein Webhook, ein Signal von einem anderen System).
  • Koordiniert mehrere Schritte mit Verzweigungen, Zusammenführungen oder Wettläufen.
  • Wird derzeit als externer Orchestrator + eine Postgres-Datenbank implementiert, wobei die meisten Arbeiten der Datenbankteil sind.

Wenn Ihre Workload aus einer einzelnen kurzen transaktionalen Anweisung besteht, benötigen Sie pg_durable nicht. Verwenden Sie eine normale INSERT / UPDATE.

So funktioniert es

Eine dauerhafte Funktion ist ein Graph von Schritten, den Sie mit einer SQL-DSL erstellen und mit df.start() übermitteln. Der Graph wird gespeichert, dann wird er von einem Hintergrundprozess ausgeführt.

Zwei Schlüsselideen:

  • Funktionsdiagramm und Ausführungszustand werden in HorizonDB selbst in den df Und duroxide Schemas gespeichert. Sicherungen, Point-in-Time-Wiederherstellung und hohe Verfügbarkeit gelten automatisch für Ihren Workflowstatus. Kein separater Orchestrator-Status muss verwaltet werden.
  • Der Hintergrundprozess wird durch shared_preload_libraries gestartet. Sie erkennt die Erweiterung nach CREATE EXTENSION und beginnt mit der Ausführung von Funktionen. Wenn die Datenbank neu gestartet wird, stellt der Worker die Verbindung zu laufenden Instanzen wieder her und setzt sie fort.

Note

Die Ausführungs-Engine in pg_durable basiert auf Duroxide, Microsofts Open-Source-Laufzeit für dauerhafte Ausführung für Rust (inspiriert durch das Durable Task Framework und Temporal). Der Name des Schemas duroxide spiegelt dies wider: Dort speichert Duroxide den Orchestrierungsverlauf, Korrelations-IDs und den Replayzustand. Die Garantien hinsichtlich deterministischer Wiedergabe, korrelierter Ereignis-IDs und dauerhafter Timer, die Sie von pg_durable erhalten, stammen direkt von Duroxide.

Aktivieren von pg_durable

Um pg_durable für Azure HorizonDB zu aktivieren, konfigurieren Sie zuerst eine Parametergruppe, und erstellen Sie dann die Erweiterung in jeder Datenbank.

Verwenden Sie die folgenden Setupartikel:

  1. Erstellen Sie eine Parametergruppe für Ihren Server.
  2. Legen Sie shared_preload_libraries so fest, dass pg_durable einbezogen wird.
  3. Legen Sie azure.extensions so fest, dass pg_durable einbezogen wird.
  4. Wenden Sie die Parametergruppe auf den Server an.
  5. Stellen Sie eine Verbindung mit jeder Zieldatenbank her, und führen Sie Folgendes aus:

Erstellen Sie die Erweiterung in jeder Datenbank, in der Sie sie verwenden möchten:

CREATE EXTENSION IF NOT EXISTS pg_durable;

CREATE EXTENSION stellt das df Schema (Funktionsdiagramme und Überwachungsansichten) und das duroxide Schema (Ausführungszustand) bereit. Der Hintergrundmitarbeiter erkennt die Erweiterung innerhalb weniger Sekunden und kann Funktionen ausführen.

Ihre erste dauerhafte Funktion

-- Start a one-step durable function
SELECT df.start('SELECT ''Hello, durable world!''');
-- Returns an 8-character instance ID, for example: a1b2c3d4

-- Check status
SELECT df.status('a1b2c3d4');

-- Get the result
SELECT df.result('a1b2c3d4');

Selbst eine Funktion mit nur einem Schritt ist robust: Wenn die Datenbank nach df.start() und bevor der Worker sie aufgreift neu gestartet wird, wird die Funktion trotzdem ausgeführt.

Note

df.start() sendet asynchron einen Workflow und gibt sofort zurück. Verwenden Sie df.list_instances(), df.instance_info(), df.status() oder df.result() für mehrstufige Workflows, um den Abschluss zu bestätigen, bevor Sie die Nebenwirkungen validieren.

Programmmodell

Eine dauerhafte Funktion ist ein Diagramm, das aus Schritten, Operatoren und integrierten Funktionen erstellt wird. Einfache SQL-Strings werden automatisch umschlossen, sodass Sie df.sql() nicht explizit aufrufen müssen.

Betriebspersonal

Operator Bedeutung Example
~> Sequenz - zuerst links, dann rechts ausführen 'SELECT 1' ~> 'SELECT 2'
& Zusammenführen – parallel ausführen, auf alle warten 'SELECT 1' & 'SELECT 2'
| Rennen - Parallel laufen, erste Siege fast_query | df.sleep(30)
?> !> If/else – Verzweigung bei einer booleschen Bedingung cond ?> then_branch !> else_branch
@> Schleife – endlos wiederholen (Präfixoperator) @> body
|=> Name - Ergebnis eines Schritts erfassen 'SELECT id FROM users LIMIT 1' |=> 'user_id'

Nützliche integrierte Funktionen

Funktion Purpose
df.sleep(seconds) Pause für N Sekunden. Dauerhaft bei Neustarts.
df.wait_for_schedule(cron) Warten Sie, bis ein Cron-Ausdruck das nächste Mal übereinstimmt.
df.wait_for_signal(name, timeout) Blockieren, bis ein externes df.signal() eintritt.
df.http(url, method, body, headers, timeout) Führen Sie einen HTTP-Aufruf als dauerhafte Aktivität durch, und versuchen Sie es bei vorübergehendem Fehler.
df.if(cond, then, else) Bedingte Verzweigung.
df.loop(body, cond) Wiederholen Sie den Vorgang, während eine SQL-Bedingung wahrheitsgetreu ist.
df.join(a, b) / df.race(a, b) Parallele und Raceausführung.
df.join3(a, b, c) Für die dreiwegige parallele Ausführung.
df.start(body, label, database) Übermitteln Sie eine Durable Function, und geben Sie deren Instanz-ID zurück.
df.cancel(id, reason) Abbrechen einer ausgeführten Instanz.
df.status(id) / df.result(id) Ergebnis prüfen.
df.explain(input) Rendern des Funktionsdiagramms für die Visualisierung.

Weitere Informationen zu allen pg_durable Features.

Variables

|=> erfasst das Ergebnis eines Schritts mit einem Namen; in späteren Schritten wird darauf verwiesen als $name.

SELECT df.start(
    'SELECT 100 AS amount' |=> 'total'
    ~> 'SELECT $total * 2 AS doubled'
);

Verwendungsbeispiele

Mehrstufige ETL-Prozesse mit Wiederholungsversuchen

Eine tägliche ETL, die bereinigt, lädt, indiziert und protokolliert:

SELECT df.start(
    'DELETE FROM target WHERE loaded_at < now() - INTERVAL ''1 day'''
    ~> 'INSERT INTO target SELECT * FROM staging'
    ~> 'REINDEX TABLE target'
    ~> 'INSERT INTO etl_log (job, finished_at) VALUES (''nightly'', now())',
    'nightly-etl'
);

Wenn die Datenbank zwischen dem DELETE und dem INSERT neu gestartet wird, setzt der Worker ab dem INSERT fort – das DELETE wird nicht erneut ausgeführt.

Geplanter Auftrag (cron)

Führen Sie jeden Wochentag um 9 Uhr einen Wartungsvorgang aus:

SELECT df.start(
    @> (
        df.wait_for_schedule('0 9 * * 1-5')
        ~> 'CALL refresh_materialized_views()'
    ),
    'weekday-refresh'
);

Wenn Sie diesen Auftrag beenden möchten, können Sie die cancel Funktion ausführen.

SELECT df.cancel('a1b2c3d4', 'stop test cron job');

Genehmigungsworkflow mit Timeout

Warten Sie bis zu 24 Stunden auf ein externes Genehmigungssignal, und führen Sie dann einen Commit oder Ablehnung durch:

SELECT df.start(
    'SELECT order_id, total FROM orders WHERE id = 1' |=> 'order'
    ~> df.wait_for_signal('approval', 86400) |=> 'sig'
    ~> df.if(
        'SELECT NOT ($sig::jsonb->>''timed_out'')::boolean
            AND ($sig::jsonb->''data''->>''approved'')::boolean',
        'UPDATE orders SET status = ''approved'' WHERE id = $order_id',
        'UPDATE orders SET status = ''rejected'' WHERE id = $order_id'
    ),
    'order-approval'
);

-- Later, approve from anywhere
SELECT df.signal('a1b2c3d4', 'approval',
                 '{"approved": true, "approver": "jane@contoso.com"}');

Dauerhafter HTTP-Aufruf

df.http() führt externe Aufrufe als dauerhafte Aktivitäten aus – 5xx-Antworten, Netzwerkfehler und Timeouts werden automatisch erneut versucht.

SELECT df.start(
    df.http('https://api.example.com/users/123', 'GET') |=> 'user'
    ~> 'INSERT INTO users_cache (data) VALUES (($user::jsonb->>''body'')::jsonb)',
    'fetch-user'
);

Weitere Informationen zur zulässigen HTTP-Sicherheit in pg_durable.

Beobachten und Bedienen

Alles ist aus SQL abfragbar. Es gibt keine separate Benutzeroberfläche oder keinen separaten Dienst, um zu lernen.

-- All instances
SELECT * FROM df.list_instances();

-- Filter by status
SELECT * FROM df.list_instances() WHERE status = 'Running';
SELECT * FROM df.list_instances() WHERE status = 'Failed';

-- Detail for one instance
SELECT * FROM df.instance_info('a1b2c3d4');

-- Execution history (useful for retried or looped functions)
SELECT * FROM df.instance_executions('a1b2c3d4', 20);

-- The function graph as it ran
SELECT * FROM df.instance_nodes('a1b2c3d4');

-- System-wide metrics
SELECT * FROM df.metrics();

So überprüfen Sie, ob der Worker aktiv ist:

SELECT epoch_id, last_seen_at, now() - last_seen_at AS time_since_last_heartbeat
FROM df._worker_epoch;

Eine time_since_last_heartbeat unter 15 Sekunden bedeutet, dass der Arbeitnehmer gesund ist. Alles, was größer oder gar keine Zeilen ist, bedeutet, dass der Worker ausgefallen ist oder nicht initialisiert wurde.

Überwachen von Workflows in Visual Studio Code

Die PostgreSQL-Erweiterung für Visual Studio Code enthält in der Ansicht Pipelines & Workflows eine Registerkarte Workflows, in der Sie pg_durableWorkflowinstanzen einsehen und den Ausführungsstatus direkt im Editor überwachen können.

Öffnen des Bereichs "Workflows"

  1. Öffnen Sie in Visual Studio Code die PostgreSQL-Erweiterung.
  2. Klicken Sie in Objekt-Explorer mit der rechten Maustaste auf Ihre Datenbank.
  3. Wählen Sie Pipelines und Workflows aus.
  4. Wählen Sie die Registerkarte "Workflows " aus.

Im linken Bereich sind PG Durable Runs aufgeführt, und im mittleren Bereich werden Details für die ausgewählte Workflowinstanz angezeigt.

Screenshot der Registerkarte

Workflow-Ausführungen überprüfen

Wenn Sie einen Workflowlauf auswählen, prüfen Sie die Zusammenfassung, um Folgendes zu bestätigen:

  • Status: completed, , runningoder failed.
  • Run ID: Eindeutiger Bezeichner für die Instanz.
  • Gestartete Zeit und Dauer: Nachverfolgen des Ausführungsfortschritts und der Leistung.
  • Detailbereich: Zusätzliche Ausführungsmetadaten.

Verwenden Sie die verfügbaren Registerkarten, um tiefer zu tauchen:

  • Graph: Visuelle Schritt-für-Schritt-Ausführungsansicht mit der Workflowstruktur und dem Schrittfluss.
  • Zeitverhalten: Zeitbasierte Ansicht für die Performanceanalyse und die Identifizierung von Engpässen.
  • Ergebnisse: Ausgabe- und ergebnisorientierte Details aus der Workflowausführung.

Bei Workflows, die sich auf KI-Pipelines beziehen, können Sie mit einer Ansicht-Pipelinedefinitionsaktion (sofern verfügbar) eine Verknüpfung von einem Workflow mit der Pipelinedefinition herstellen, die nützlich ist, um das Verhalten bei ausführungsübergreifenden Vorgängen zu vergleichen oder Regressionen zu untersuchen.

Identität und Isolation

Dauerhafte Funktionen werden mit den Berechtigungen des Benutzers ausgeführt, der sie übermittelt hat, nicht mit den Berechtigungen des Mitarbeiters. pg_durable erfasst zum Zeitpunkt der Übermittlung sowohl session_user als auch current_user, sodass Funktionen, die im Kontext von SET ROLE übermittelt werden, mit dieser effektiven Rolle ausgeführt werden.

Dies bedeutet Folgendes:

  • Benutzer sehen und ändern nur Daten, auf die sie bereits über Berechtigungen für den Zugriff verfügen.
  • Benutzer ohne Superuser-Rechte können ihre Rechte nicht ausweiten, indem sie eine persistente Funktion übermitteln.
  • Mehrinstanzenfähige Workloads bleiben isoliert, solange Ihr Rollen- und Zuweisungsmodell korrekt ist.

Interaktion mit Replikaten, Sicherungen und PITR

  • Sicherung und PITR. Funktionsdiagramm (df Schema) und Ausführungszustand (duroxide Schema) werden in regulären Tabellen gespeichert und sind in HorizonDB-Sicherungen enthalten. Eine Zeitpunktwiederherstellung stellt beide wieder her.
  • Lesereplikate. Der Hintergrundworker wird nur auf der primären Instanz ausgeführt. Lesereplikate können die df.* Überwachungsansichten abfragen, aber keine Funktionen ausführen.
  • Failover. Nach einem Failover setzt der Worker auf der neuen Primärinstanz dort fort, wo die bisherige Primärinstanz aufgehört hat. Ausführungsinstanzen setzen die Ausführung ab ihrem letzten Prüfpunkt fort.

Im Vergleich zu externen Orchestratoren

Aspect Externer Orchestrator pg_durable
Einsatz Separater Dienst, separate Identität, separater Zustandsspeicher Eine Datenbank
Zustandsbeständigkeit Speicherschicht des Orchestrators Dieselben Sicherungen, HA und PITR wie Ihre Daten
Identity Worker werden unter einer Dienstidentität ausgeführt Funktionen werden als übermittelnder Benutzer ausgeführt
Fehlerarten Netzwerk zwischen Orchestrator und Datenbank Keine - gleicher Prozess
Am besten geeignet für: Systemübergreifende Orchestrierung, die viele Dienste berührt Arbeitsauslastungen, bei denen sich der Großteil der Arbeit in oder in der Nähe von Postgres befindet

pg_durable versucht nicht, externe Orchestratoren für systemübergreifende Pipelines zu ersetzen. Dies ist die richtige Wahl, wenn die meiste Arbeit Datenbankarbeit ist – Einbettungen, Transformationen, KI-Aufrufe, geplante Wartung und Hinzufügen eines weiteren Diensts ist mehr Kosten als Nutzen.

Einschränkungen während der Vorschau

  • df.http() Wiederholungsversuche bei 5xx-Fehlern und Netzwerkfehlern. 4xx-Antworten werden an den Workflow zurückgegeben, damit Sie sie selbst behandeln; sie werden nicht automatisch erneut versucht.
  • Der Hintergrundarbeitsdienst bedient pro Instanz eine einzelne Datenbank. Mehrfach-Datenbank-Fan-Out wird über df.start(..., database => 'other_db') aus einer Funktion unterstützt, die in der Datenbank des Workers ausgeführt wird.
  • Funktionsdefinitionen und Ausführungsstatus sind während pg_durableder Vorschau nicht zwischen Hauptversionen portierbar. Löschen oder Abbrechen der Ausführung von Instanzen vor dem Upgrade.