Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
pg_durable je trvalý prováděcí modul uvnitř Azure HorizonDB. Umožňuje definovat dlouhotrvající vícestupňové pracovní postupy SQL (vkládání kanálů, úloh ETL, volání AI, naplánované úlohy, toky schválení) a spouštět je se stejnými zárukami spolehlivosti, které byste očekávali od vyhrazeného orchestrátoru, jako je Durable Functions, aniž byste opustili Postgres.
pg_durable je také vrstvou pro spouštění, která tvoří základ odolných AI pipeline. Pokud používáte AI pipeline, právě pg_durable zajišťuje, že přežijí pády, po selhání se znovu spustí a pokračují od posledního dokončeného kroku.
Note
pg_durable je ve verzi Preview.
Co znamená "trvalý"
Trvalá funkce v pg_durable se v každém kroku ukládá na disk. Získáte tak konkrétní sadu záruk, které nedostanete z prostého BEGIN ... COMMIT bloku nebo úlohy cron:
- Přežije chybové ukončení a restartování databáze. Dokončené kroky se po opětovném spuštění serveru znovu nespustí. Probíhající kroky se obnoví z posledního kontrolního bodu. Čekající kroky se spustí, když se pracovní proces vrátí do režimu online.
- Vydrží i dlouhé čekání. Workflow může být uspáno na celé hodiny, čekat na plán cronu nebo na externí signál a přesto pokračovat tam, kde skončilo.
- Přežije selhání. Neúspěšné kroky je možné opakovat automaticky bez opětovného spuštění celé funkce.
- Zaznamenává identitu. Funkce se spustí s oprávněními uživatele, který ho spustil, nikoli s oprávněními pracovního procesu na pozadí. Multitenantní úlohy zůstávají izolované.
- Zůstává pozorovatelný z SQL. Stav, historii, počet spuštění a výstupy můžete zkontrolovat pomocí stejného rozhraní, které používáte pro všechno ostatní v HorizonDB:
SELECTpříkaz.
Co odolnost automaticky nezajistí: sama o sobě nezajistí, že bude bezpečné opakovat neidempotentní externí operace. Pokud krok volá externí rozhraní API, které účtuje peníze, navrhujte krok tak, aby byl idempotentní (například předáním klíče idempotence).
Kdy použít pg_durable
Použijte pg_durable, když s tím musíte pracovat:
- Trvá to dost dlouho na to, aby to selhalo až uprostřed (generování embeddingů nad miliony řádků, vícekroková úloha ETL, backfill).
- V případě selhání je potřeba akci zopakovat, aniž by se znovu prováděly části, které již byly úspěšně dokončeny.
- Musí běžet podle harmonogramu (každou hodinu, každý pracovní den v 9:00).
- Potřebuje počkat na externí událost (schválení, webhook, signál z jiného systému).
- Koordinuje více kroků pomocí větvení, spojení nebo závodu.
- V současné době se implementuje jako externí orchestrátor + databáze Postgres, kde většina práce je součástí databáze.
Pokud je vaše úloha jedním krátkým transakčním příkazem, nepotřebujete pg_durable. Používejte běžnou INSERT / UPDATE.
Jak to funguje
Odolná funkce je graf kroků, které vytvoříte pomocí SQL DSL a odešlete s df.start(). Graf je trvalý a pak ho spustí pracovní proces na pozadí.
Dvě klíčové myšlenky:
-
Graf funkcí a stav provádění se ukládají do samotné databáze HorizonDB ve
dfschématech aduroxideschématech. Zálohy, obnovení k určitému bodu v čase a vysoká dostupnost se automaticky vztahují na stav vašeho workflow. Není potřeba spravovat samostatný stav orchestrátoru. - Proces na pozadí je spuštěn pomocí
shared_preload_libraries. PoCREATE EXTENSIONzjistí rozšíření a začne vykonávat funkce. Pokud se databáze restartuje, pracovní proces se znovu připojí ke spuštěným instancím a obnoví je.
Note
Vykonávací modul uvnitř pg_durable je postavený na Duroxide, open-source běhovém prostředí společnosti Microsoft pro odolné vykonávání v Rustu (inspirovaném frameworkem Durable Task Framework a platformou Temporal). Název schématu duroxide to odráží: právě tam Duroxide ukládá historii orchestrace, identifikátory korelace a stav přehrávání. Záruky deterministického přehrávání, ID korelovaných událostí a odolných časovačů, které získáte od pg_durable, pocházejí přímo z Duroxide.
Povolit pg_durable
Pokud chcete povolit pg_durable v Azure HorizonDB, nejprve nakonfigurujte skupinu parametrů a pak vytvořte rozšíření v každé databázi.
Použijte tyto články o nastavení:
- Vytvořte skupinu parametrů pro váš server.
- Nastavte
shared_preload_librariestak, aby zahrnovalpg_durable. - Nastavte
azure.extensionstak, aby zahrnovalpg_durable. - Přiřaďte skupinu parametrů k serveru.
- Připojte se ke každé cílové databázi a spusťte:
Vytvořte rozšíření v každé databázi, ve které ho chcete použít:
CREATE EXTENSION IF NOT EXISTS pg_durable;
CREATE EXTENSION zajišťuje schéma df (grafy funkcí a monitorovací zobrazení) a schéma duroxide (stav běhu). Pracovní proces na pozadí zjistí rozšíření během několika sekund a je připravený ke spuštění funkcí.
Vaše první odolná funkce
-- 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');
Dokonce i jednokroková funkce je robustní: pokud se databáze restartuje po df.start() a předtím, než ji worker vyzvedne, funkce se stále spustí.
Note
df.start() odešle pracovní postup asynchronně a vrátí okamžitě. U vícekrokových pracovních postupů použijte df.list_instances(), df.instance_info(), df.status(), nebo df.result() potvrďte dokončení před ověřením vedlejších účinků.
Programový model
Odolná funkce je graf vytvořený z kroků, operátorů a předdefinovaných funkcí. Prosté řetězce SQL jsou automaticky přepsány, takže nemusíte volat df.sql() explicitně.
Operators
| Operator | Význam | Example |
|---|---|---|
~> |
Sekvence – běh vlevo, poté vpravo | 'SELECT 1' ~> 'SELECT 2' |
& |
Připojení – paralelní spuštění, čekání na vše | 'SELECT 1' & 'SELECT 2' |
| |
Závod - běžet paralelně, první vyhrává | fast_query | df.sleep(30) |
?>
!>
|
If / else – větvení podle booleovské podmínky | cond ?> then_branch !> else_branch |
@> |
Smyčka – opakovat donekonečna (prefixový operátor) | @> body |
|=> |
Název – zachycení výsledku kroku | 'SELECT id FROM users LIMIT 1' |=> 'user_id' |
Užitečné vestavěné funkce
| Function | Purpose |
|---|---|
df.sleep(seconds) |
Pozastavit po dobu N sekund. Přetrvá i po restartu. |
df.wait_for_schedule(cron) |
Počkejte, až se cron výraz příště vyhodnotí jako shodný. |
df.wait_for_signal(name, timeout) |
Zablokujte, dokud nedorazí externí df.signal() . |
df.http(url, method, body, headers, timeout) |
Proveďte volání HTTP jako odolnou aktivitu s opakováním při přechodném selhání. |
df.if(cond, then, else) |
Podmíněná větev |
df.loop(body, cond) |
Opakujte, když je podmínka SQL pravdivá. |
df.join(a, b) / df.race(a, b) |
Paralelní a souběžné spouštění. |
df.join3(a, b, c) |
Pro trojcestné paralelní spouštění. |
df.start(body, label, database) |
Odešlete trvalou funkci a vraťte ID její instance. |
df.cancel(id, reason) |
Zrušení spuštěné instance |
df.status(id) / df.result(id) |
Zkontrolujte výsledek. |
df.explain(input) |
Vykreslení grafu funkce pro vizualizaci |
Přečtěte si další informace o všech funkcích pg_durable.
Variables
|=> uloží výsledek kroku pod názvem; pozdější kroky na něj odkazují jako na $name.
SELECT df.start(
'SELECT 100 AS amount' |=> 'total'
~> 'SELECT $total * 2 AS doubled'
);
Příklady použití
Vícekrokové ETL s opakovanými pokusy
Denní ETL, která vyčistí, načte, zaindexuje a zaznamená:
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'
);
Pokud se databáze restartuje mezi DELETE a INSERT, pracovní proces pokračuje od INSERT – znovu nespustí DELETE.
Naplánovaná úloha (cron)
Spusťte úlohu údržby každý pracovní den v 9:00:
SELECT df.start(
@> (
df.wait_for_schedule('0 9 * * 1-5')
~> 'CALL refresh_materialized_views()'
),
'weekday-refresh'
);
Pokud chcete tuto úlohu zastavit, můžete funkci spustit cancel .
SELECT df.cancel('a1b2c3d4', 'stop test cron job');
Proces schvalování s časovým limitem
Počkejte až 24 hodin na signál externího schválení a pak potvrďte nebo odmítněte:
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"}');
Trvalé volání HTTP
df.http() provádí externí volání jako trvalé aktivity – automaticky se opakují odpovědi 5xx, chyby sítě a vypršení časových limitů.
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'
);
Přečtěte si další informace o povoleném zabezpečení HTTP v pg_durable.
Sledování a ovládání
Vše je možné dotazovat z SQL. Není potřeba učit se používat žádné samostatné uživatelské rozhraní ani službu.
-- 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();
Chcete-li ověřit, že je worker aktivní:
SELECT epoch_id, last_seen_at, now() - last_seen_at AS time_since_last_heartbeat
FROM df._worker_epoch;
Hodnota time_since_last_heartbeat kratší než 15 sekund znamená, že worker je zdravý. Cokoli větší nebo vůbec žádné řádky znamená, že pracovní proces je mimo provoz nebo není inicializován.
Monitorování pracovních postupů v Visual Studio Code
Rozšíření PostgreSQL pro Visual Studio Code zahrnuje kartu Pracovní postupy v zobrazení Pipelines & Workflows, kde můžete prohlížet pg_durable instance pracovních postupů a sledovat stav běhu přímo z editoru.
Otevření podokna Pracovní postupy
- V Visual Studio Code otevřete rozšíření PostgreSQL.
- V Průzkumník objektů klikněte pravým tlačítkem na databázi.
- Vyberte Kanály a pracovní postupy.
- Vyberte kartu Pracovní postupy .
V levém podokně je uveden seznam PG Durable Runs a prostřední podokno zobrazuje podrobnosti o vybrané instanci pracovního postupu.
Zkontrolovat běhy pracovního postupu
Když vyberete spuštění pracovního postupu, zkontrolujte souhrn a ověřte následující:
-
Stav:
completed,runningnebofailed. - ID spuštění: Jedinečný identifikátor instance.
- Čas zahájení a doba trvání: Sledujte průběh spouštění a výkon.
- Panel podrobností: Další metadata o spuštění.
Použijte dostupné karty k podrobnějšímu prozkoumání:
- Graf: Vizuální podrobné zobrazení provádění znázorňující strukturu pracovního postupu a tok kroku
- Časování: Zobrazení zaměřené na dobu trvání pro analýzu výkonu a identifikaci kritických bodů
- Výsledky: Výstup a podrobnosti orientované na výsledky z provádění pracovního postupu.
U pracovních postupů souvisejících s pipeline AI umožňuje akce Zobrazit definici pipeline (je-li k dispozici) přejít z běhu pracovního postupu zpět k definici příslušné pipeline, což je užitečné při porovnávání chování napříč běhy nebo při zkoumání regresí.
Identita a izolace
Odolné funkce se spouštějí s oprávněními uživatele, který je odeslal, ne s oprávněními pracovního procesu.
pg_durable při odeslání zachytí jak session_user, tak current_user, takže funkce odeslané v kontextu SET ROLE se spouštějí s touto efektivní rolí.
To znamená:
- Uživatelé vidí a upravují data, ke kterým už mají oprávnění pro přístup.
- Uživatelé bez oprávnění superuživatele nemohou získat vyšší oprávnění odesláním trvalé funkce.
- Multitenantní úlohy zůstávají izolované, pokud je váš model rolí a oprávnění správný.
Interakce s replikami, zálohami a PITR
- Zálohování a PITR. Graf funkcí (schéma) a stav spouštění (
dfduroxideschéma) se ukládají v pravidelných tabulkách a jsou zahrnuty do záloh HorizonDB. Obnovení k určitému časovému bodu obnoví oboje. - Repliky pro čtení Pracovní proces na pozadí běží jenom na primárním serveru. Repliky pro čtení se můžou dotazovat na
df.*zobrazení monitorování, ale nespouštějí funkce. - Převzetí služeb při selhání. Po převzetí služeb při selhání pracovní proces na novém primárním serveru převezme místo, kde starý primární proces skončil. Běžící instance pokračují od posledního kontrolního bodu.
Porovnání s externími orchestrátory
| Aspect | Externí orchestrátor | pg_durable |
|---|---|---|
| Deployment | Samostatná služba, samostatná identita, samostatné úložiště stavů | Jedna databáze |
| Odolnost stavu | Vrstva úložiště nástroje Orchestrator | Stejné zálohování, vysoká dostupnost a obnovení k bodu v čase jako u vašich dat |
| Identity | Pracovní procesy běží pod identitou služby | Funkce se spouštějí jako odesílající uživatel. |
| Režimy selhání | Síť mezi orchestrátorem a databází | Žádné – stejný proces |
| Nejlepší pro | Orchestrace mezi systémy, která se dotýká mnoha služeb | Úlohy, ve kterých je většina práce v Postgresu nebo blízko |
pg_durable nepokouší se nahradit externí orchestrátory pro kanály mezi systémy. Je to správná volba, když je většina práce v databázi – vkládání, transformace, volání AI, plánovaná údržba – a přidání další služby je nákladnější než výhoda.
Omezení ve verzi Preview
-
df.http()opakování požadavku při chybách 5xx a síťových chybách. Odpovědi 4xx se vrátí do pracovního postupu, který budete zpracovávat; nebudou se opakovat automaticky. - Služba běžící na pozadí obsluhuje v rámci jedné instance jednu databázi. Rozvětvení do více databází je podporováno pomocí
df.start(..., database => 'other_db')z funkce spuštěné v databázi pracovního procesu. - Definice funkcí a stav provádění nejsou mezi hlavními verzemi
pg_durablev Preview přenosné. Vyprázdnění nebo zrušení spuštěných instancí před upgradem