Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
pg_durable è il motore di esecuzione durevole all'interno di Azure HorizonDB. Consente di definire flussi di lavoro SQL a esecuzione prolungata e in più passaggi (pipeline di incorporamento, processi ETL, chiamate di intelligenza artificiale, processi pianificati, flussi di approvazione) ed eseguirli con le stesse garanzie di affidabilità previste da un agente di orchestrazione dedicato come Durable Functions, senza uscire da Postgres.
pg_durable è anche il livello di esecuzione alla base delle pipeline di IA durevoli. Se si usano pipeline di IA, pg_durable è ciò che consente loro di resistere agli arresti anomali, ritentare in caso di errore e riprendere dall'ultimo passaggio completato.
Note
pg_durable è in anteprima.
Cosa significa "durevole"
Una funzione durevole in pg_durable viene salvata in modo permanente su disco a ogni passaggio. Ciò offre un set specifico di garanzie che non si ottengono da un blocco normale BEGIN ... COMMIT o da un processo cron:
- Sopravvive agli arresti anomali e ai riavvii del database. I passaggi completati non vengono eseguiti di nuovo quando viene eseguito il backup del server. Le fasi in corso riprendono dall'ultimo punto di controllo. I passaggi in sospeso vengono eseguiti quando il lavoratore torna online.
- Resiste a lunghe attese. Un flusso di lavoro può restare sospeso per ore, attendere una pianificazione Cron o bloccarsi in attesa di un segnale esterno, e comunque riprendere da dove si era interrotto.
- Resiste agli errori. I passaggi non riusciti possono essere ritentati automaticamente senza riesecuzione dell'intera funzione.
- Rileva l'identità. Una funzione viene eseguita con i privilegi dell'utente che lo ha avviato, non con i privilegi del ruolo di lavoro in background. I carichi di lavoro multi-tenant rimangono isolati.
- Rimane monitorabile tramite SQL. È possibile esaminare lo stato, la cronologia, il numero di esecuzioni e gli output tramite la stessa interfaccia usata per tutto il resto in HorizonDB: un'istruzione
SELECT.
Cosa la durabilità non fa automaticamente: non rende automaticamente sicuro ritentare operazioni esterne non idempotenti. Se un passaggio chiama un'API esterna che addebita denaro, progettare il passaggio come idempotente, ad esempio passando una chiave di idempotenza.
Quando usare pg_durable
Usare pg_durable quando è necessario eseguire questa operazione:
- Dura abbastanza a lungo da potersi interrompere a metà (generazione di embedding su milioni di righe, un processo ETL articolato in più fasi, un backfill).
- Deve essere ritentata in caso di errore senza ripetere le parti già riuscite.
- Deve essere eseguito in base a una pianificazione (ogni ora, ogni giorno feriale alle 9:00).
- Deve attendere un evento esterno (un'approvazione, un webhook, un segnale da un altro sistema).
- Coordina più passaggi con la diramazione, l'aggiunta o la corsa.
- È attualmente implementato come agente di orchestrazione esterno e un database Postgres, dove la maggior parte del lavoro è la parte del database.
Se il carico di lavoro è una singola istruzione transazionale breve, non è necessario pg_durable. Usare un normale INSERT / UPDATE.
Come funziona
Una funzione duratura è un grafo di passaggi che si crea con una DSL SQL e si invia con df.start(). Il grafico viene salvato in modo permanente e quindi viene eseguito da un lavoratore in background.
Due idee chiave:
-
Il grafico delle funzioni e lo stato di esecuzione vengono archiviati in HorizonDB stesso, negli
dfschemi eduroxide. I backup, il ripristino temporizzato e la disponibilità elevata si applicano automaticamente allo stato del flusso di lavoro. Nessuno stato separato dell'orchestratore da gestire. - Il lavoratore in background viene avviato da
shared_preload_libraries. Rileva l'estensione dopoCREATE EXTENSIONe inizia l'esecuzione di funzioni. Se il database si riavvia, il lavoratore si ricollega alle istanze in esecuzione e ne riprende l'esecuzione.
Note
Il motore di esecuzione all'interno di pg_durable è basato su Duroxide, il runtime open source di Microsoft per l'esecuzione durevole in Rust (ispirato al Durable Task Framework e a Temporal). Il duroxide nome dello schema riflette questo aspetto: in questo caso Duroxide rende persistente la cronologia dell'orchestrazione, gli ID di correlazione e lo stato di riproduzione. Le garanzie di riproduzione deterministica, ID evento correlato e timer persistente che si ottengono da pg_durable provengono direttamente da Duroxide.
Abilitare pg_durable
Per abilitare pg_durable in Azure HorizonDB, configurare prima un gruppo di parametri e quindi creare l'estensione in ogni database.
Usare questi articoli di installazione:
- Creare un gruppo di parametri per il server.
- Impostare
shared_preload_librariesper includerepg_durable. - Impostare
azure.extensionsper includerepg_durable. - Applicare il gruppo di parametri al server.
- Connettersi a ogni database di destinazione ed eseguire:
Creare l'estensione in ogni database in cui si vuole usarla:
CREATE EXTENSION IF NOT EXISTS pg_durable;
CREATE EXTENSION predispone lo schema df (grafi di funzioni e viste di monitoraggio) e lo schema duroxide (stato di esecuzione). Il processo in background riconosce l'estensione entro pochi secondi ed è pronto a eseguire le funzioni.
La prima funzione durevole
-- 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');
Anche una funzione a singolo passaggio è affidabile: se il database viene riavviato dopo df.start() e prima che il worker la prenda in carico, la funzione verrà comunque eseguita.
Note
df.start() invia un flusso di lavoro in modo asincrono e restituisce immediatamente. Per i flussi di lavoro in più passaggi, usare df.list_instances(), df.instance_info(), df.status()o df.result() per confermare il completamento prima di convalidare gli effetti collaterali.
Modello di programma
Una funzione durevole è un grafico basato su passaggi, operatori e funzioni predefinite. Le stringhe SQL semplici vengono racchiuse automaticamente, quindi non è necessario chiamare df.sql() esplicitamente.
Operators
| Operatore | Meaning | Esempio |
|---|---|---|
~> |
Sequenza: eseguire prima a sinistra, poi a destra | 'SELECT 1' ~> 'SELECT 2' |
& |
Join - eseguire in parallelo, attendere tutto | 'SELECT 1' & 'SELECT 2' |
| |
Gara - corsa in parallelo, prima vittoria | fast_query | df.sleep(30) |
?>
!>
|
If / else - diramazione in base a una condizione booleana | cond ?> then_branch !> else_branch |
@> |
Ciclo - ripeti all'infinito (operatore prefisso) | @> body |
|=> |
Nome - acquisire il risultato di un passaggio | 'SELECT id FROM users LIMIT 1' |=> 'user_id' |
Elementi integrati utili
| Function | Purpose |
|---|---|
df.sleep(seconds) |
Sospendere per N secondi. Permane dopo i riavvii. |
df.wait_for_schedule(cron) |
Attendere fino alla successiva corrispondenza di un'espressione Cron. |
df.wait_for_signal(name, timeout) |
Attendere fino all'arrivo di un df.signal() esterno. |
df.http(url, method, body, headers, timeout) |
Effettuare una chiamata HTTP come attività durevole, con nuovi tentativi in caso di errore temporaneo. |
df.if(cond, then, else) |
Ramo condizionale. |
df.loop(body, cond) |
Ripeti mentre una condizione SQL è vera. |
df.join(a, b) / df.race(a, b) |
Esecuzione parallela e competitiva. |
df.join3(a, b, c) |
Per l'esecuzione parallela a tre vie. |
df.start(body, label, database) |
Inviare una Durable Function e restituire l'ID dell'istanza. |
df.cancel(id, reason) |
Annullare un'istanza in esecuzione. |
df.status(id) / df.result(id) |
Esaminare il risultato. |
df.explain(input) |
Visualizza il grafico della funzione. |
Altre informazioni su tutte le funzionalità di pg_durable.
Variabili
|=> acquisisce il risultato di un passaggio con un nome; i passaggi successivi lo fanno riferimento come $name.
SELECT df.start(
'SELECT 100 AS amount' |=> 'total'
~> 'SELECT $total * 2 AS doubled'
);
Esempi di utilizzo
ETL in più passaggi con nuovi tentativi
Un ETL giornaliero che pulisce, carica, indicizza e registra:
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'
);
Se il database si riavvia tra DELETE e INSERT, il lavoratore riprende da INSERT: non riesegue DELETE.
Attività pianificata (cron)
Eseguire un'attività di manutenzione ogni giorno feriale alle 9:00:
SELECT df.start(
@> (
df.wait_for_schedule('0 9 * * 1-5')
~> 'CALL refresh_materialized_views()'
),
'weekday-refresh'
);
Se si desidera arrestare questa attività, è possibile eseguire la funzione cancel.
SELECT df.cancel('a1b2c3d4', 'stop test cron job');
Flusso di lavoro di approvazione con timeout
Attendere fino a 24 ore per un segnale di approvazione esterno, quindi eseguire il commit o rifiutare:
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"}');
Chiamata HTTP persistente
df.http() effettua chiamate esterne come attività durevoli: le risposte 5xx, gli errori di rete e i timeout vengono ritentati automaticamente.
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'
);
Altre informazioni sulla sicurezza HTTP consentita in pg_durable.
Osservare e operare
Tutto può essere interrogato tramite SQL. Non esiste un'interfaccia utente o un servizio separato da apprendere.
-- 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();
Per verificare che il lavoratore sia attivo:
SELECT epoch_id, last_seen_at, now() - last_seen_at AS time_since_last_heartbeat
FROM df._worker_epoch;
Un time_since_last_heartbeat valore inferiore a 15 secondi indica che il lavoratore è integro. Qualsiasi valore superiore, oppure nessuna riga, significa che il lavoratore non è attivo o non è stato inizializzato.
Monitorare i flussi di lavoro in Visual Studio Code
L'estensione PostgreSQL per Visual Studio Code include una scheda Flussi di lavoro nella pipeline e Visualizzazione flussi di lavoro, in cui è possibile esaminare le pg_durable istanze del flusso di lavoro e monitorare lo stato di esecuzione dall'editor.
Aprire il riquadro Flussi di lavoro
- In Visual Studio Code aprire l'estensione PostgreSQL.
- In Esplora oggetti fare clic con il pulsante destro del mouse sul database.
- Selezionare Pipeline e flussi di lavoro.
- Selezionare la scheda Flussi di lavoro .
Nel riquadro sinistro sono elencati i PG Durable Runs e il riquadro centrale mostra i dettagli dell'istanza del workflow selezionata.
Esaminare le esecuzioni del flusso di lavoro
Quando selezioni un'esecuzione di un flusso di lavoro, esamina il riepilogo per verificare:
-
Stato:
completed,runningofailed. - ID dell'esecuzione: identificatore univoco per l'istanza.
- Tempo e durata avviati: tenere traccia dello stato e delle prestazioni dell'esecuzione.
- Pannello dei dettagli: metadati di esecuzione aggiuntivi.
Usa le schede disponibili per approfondire:
- Grafo: vista dell'esecuzione passo dopo passo che mostra la struttura del flusso di lavoro e il flusso dei passaggi.
- Intervallo: visualizzazione incentrata sulla durata per l'analisi delle prestazioni e l'identificazione dei colli di bottiglia.
- Risultati: dati di output e dettagli relativi ai risultati dell'esecuzione del flusso di lavoro.
Per i flussi di lavoro correlati alle pipeline di intelligenza artificiale, un'azione View pipeline definition (se disponibile) consente di collegare un flusso di lavoro alla definizione della pipeline, utile per confrontare il comportamento tra esecuzioni o analizzare le regressioni.
Identità e isolamento
Le funzioni permanenti vengono eseguite con i privilegi dell'utente che li ha inviati, non con i privilegi del ruolo di lavoro.
pg_durable acquisisce sia session_user sia current_user al momento dell'invio, quindi le funzioni inviate nel contesto di SET ROLE vengono eseguite con quel ruolo effettivo.
Ciò significa:
- Gli utenti vedono e modificano solo i dati a cui dispongono già delle autorizzazioni per l'accesso.
- Gli utenti non privilegiati non possono elevare i privilegi inviando una funzione Durable.
- I carichi di lavoro multi-tenant rimangono isolati purché il ruolo e il modello di concessione sono corretti.
Interazione con repliche, backup e PITR
- Backup e PITR. Il grafico delle funzioni (
dfschema) e lo stato di esecuzione (duroxideschema) vengono archiviati in tabelle regolari e sono inclusi nei backup di HorizonDB. Un ripristino a un punto nel tempo ripristina entrambi. - Repliche in lettura. Il lavoratore in background viene eseguito solo nel database primario. Le repliche di lettura possono interrogare le
df.*viste di monitoraggio, ma non possono eseguire funzioni. - Failover. Dopo un failover, il lavoratore nel nuovo database primario rileva dove è stata interrotta la precedente replica primaria. Le istanze attive riprendono dall'ultimo checkpoint.
Rispetto agli agenti di orchestrazione esterni
| Aspect | Agente di orchestrazione esterno | pg_durable |
|---|---|---|
| Distribuzione | Servizio separato, identità separata, archivio dello stato separato | Un database |
| Durabilità dello stato | Livello di archiviazione di Orchestrator | Gli stessi backup, HA e PITR dei tuoi dati |
| Identity | I lavoratori vengono eseguiti con un'identità del servizio | Le funzioni vengono eseguite come l'utente che ha effettuato l'invio |
| Modalità di errore | Rete tra agente di orchestrazione e database | Nessuno - stesso processo |
| Migliore per | Orchestrazione tra sistemi che tocca molti servizi | Carichi di lavoro in cui la maggior parte del lavoro si trova in o vicino a Postgres |
pg_durable non sta tentando di sostituire orchestratori esterni per le pipeline tra sistemi. È la scelta giusta quando la maggior parte del lavoro è il lavoro del database, ovvero incorporamenti, trasformazioni, chiamate di intelligenza artificiale, manutenzione pianificata e l'aggiunta di un altro servizio è più conveniente del vantaggio.
Limitazioni durante l'anteprima
-
df.http()riprova in caso di errori 5xx e di rete. Le risposte 4xx vengono restituite al flusso di lavoro affinché tu le gestisca; non vengono ritentate automaticamente. - Il lavoratore in background gestisce un solo database per istanza. Il fan-out multi-database è supportato tramite
df.start(..., database => 'other_db')da una funzione in esecuzione nel database del lavoratore. - Le definizioni di funzione e lo stato di esecuzione non sono portabili tra le versioni principali di
pg_durabledurante l'anteprima. Svuotare o annullare le istanze in esecuzione prima dell'aggiornamento.
Contenuti correlati
- Implementare pipeline di intelligenza artificiale resilienti in Azure HorizonDB (Anteprima)
- funzioni AI nell'estensione azure_ai per Azure HorizonDB (anteprima)
- Generare incorporamenti vettoriali usando la funzione di intelligenza artificiale create_embeddings() (anteprima)
- Consentire le estensioni in Azure HorizonDB (anteprima)
- Duroxide in GitHub