Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
pg_durable é o mecanismo de execução durável dentro do Azure HorizonDB. Ele permite que você defina fluxos de trabalho sql de várias etapas de execução longa (incorporando pipelines, trabalhos ETL, chamadas de IA, trabalhos agendados, fluxos de aprovação) e execute-os com as mesmas garantias de confiabilidade que você esperaria de um orquestrador dedicado como Durable Functions, sem sair do Postgres.
pg_durable também é a camada de execução subjacente a pipelines duráveis de IA. Se você estiver usando pipelines de IA, pg_durable é o que permite que eles sobrevivam a falhas, sejam tentados novamente em caso de falha e retomem a partir da última etapa concluída.
Note
pg_durable está em versão prévia.
O que significa "durável"
Uma função durável em pg_durable é mantida em disco a cada etapa. Isso oferece a você um conjunto específico de garantias que você não obtém com um bloco BEGIN ... COMMIT simples nem com uma tarefa cron:
- Sobrevive a falhas e reinicializações do banco de dados. As etapas concluídas não são executadas novamente quando o servidor volta a ficar online. As etapas em andamento são retomadas a partir do último checkpoint. As etapas pendentes são executadas quando o trabalho fica online novamente.
- Resiste a longos períodos de espera. Um fluxo de trabalho pode dormir por horas, aguardar uma agenda cron ou bloquear um sinal externo e ainda continuar de onde parou.
- Resiste a falhas. As etapas com falha podem ser repetidas automaticamente sem executar novamente toda a função.
- Captura a identidade. Uma função é executada com os privilégios do usuário que a iniciou, não os privilégios do trabalho em segundo plano. As cargas de trabalho multitenant permanecem isoladas.
- Continua observável no SQL. Você pode inspecionar o status, o histórico, a contagem de execuções e as saídas por meio da mesma interface usada para todo o resto no HorizonDB: uma
SELECTinstrução.
O que a durabilidade não faz automaticamente: não torna seguras, por si só, as operações externas não idempotentes para nova tentativa. Se uma etapa chamar uma API externa que gera cobrança, projete a etapa de modo que ela seja idempotente (por exemplo, passando uma chave de idempotência).
Quando usar pg_durable
Use pg_durable quando tiver que trabalhar nisso:
- Demora o suficiente para falhar no meio (geração de incorporações em milhões de linhas, um trabalho ETL multi-etapa, um provisionamento).
- Precisa ser tentado novamente se houver falha, sem refazer as partes que já foram concluídas com sucesso.
- Precisa ser executado conforme uma programação (a cada hora, todos os dias úteis às 9h).
- Precisa aguardar um evento externo (uma aprovação, um webhook, um sinal de outro sistema).
- Coordena várias etapas com ramificação, junção ou corrida.
- Atualmente, é implementado como um orquestrador externo + um banco de dados Postgres, onde a maior parte do trabalho é a parte do banco de dados.
Se sua carga de trabalho for uma única instrução transacional curta, você não precisará pg_durable. Use um comum INSERT / UPDATE.
Como funciona
Uma função durável é um grafo de etapas criado por você com uma DSL SQL e enviado com df.start(). O grafo é mantido e, em seguida, um trabalho em segundo plano o executa.
Duas ideias importantes:
-
O grafo de funções e o estado de execução são armazenados no próprio HorizonDB, nos esquemas
dfeduroxide. Backups, restauração pontual e alta disponibilidade se aplicam automaticamente ao seu estado de fluxo de trabalho. Nenhum estado de orquestrador à parte a ser gerenciado. - O trabalho em segundo plano é iniciado por
shared_preload_libraries. Ele detecta a extensão apósCREATE EXTENSIONe começa a executar funções. Se o banco de dados é reiniciado, o trabalho se reanexa a instâncias em execução e as retoma.
Note
O mecanismo de execução dentro de pg_durable é baseado em Duroxide, o runtime open-source de execução durável da Microsoft para Rust (inspirado no Durable Task Framework e no Temporal). O nome do esquema duroxide reflete isso: é nele que Duroxide mantém o histórico de orquestração, as IDs de correlação e o estado de reprodução. As garantias de reprodução determinística, ID de evento correlacionada e temporizador durável obtidas por você de pg_durable vêm diretamente do Duroxide.
Habilitar pg_durable
Para habilitar pg_durable no Azure HorizonDB, primeiro configure um grupo de parâmetros e crie a extensão em cada banco de dados.
Use estes artigos de instalação:
- Crie um grupo de parâmetros para o servidor.
- Defina
shared_preload_librariespara incluirpg_durable. - Defina
azure.extensionspara incluirpg_durable. - Aplique o grupo de parâmetros ao servidor.
- Conecte-se a cada banco de dados de destino e execute:
Crie a extensão em cada banco de dados em que você deseja usá-la:
CREATE EXTENSION IF NOT EXISTS pg_durable;
CREATE EXTENSION provisiona o df esquema (grafos de função e exibições de monitoramento) e o duroxide esquema (estado de execução). O trabalho em segundo plano detecta a extensão dentro de alguns segundos e está pronto para executar funções.
Sua primeira função durável
-- 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');
Até mesmo uma função de etapa única é durável: se o banco de dados reiniciar depois de df.start() e antes do trabalho a processar, a função será executada mesmo assim.
Note
df.start() envia um fluxo de trabalho de forma assíncrona e retorna imediatamente. Para fluxos de trabalho de várias etapas, use df.list_instances(), df.instance_info(), df.status() ou df.result() para confirmar a conclusão antes de validar os efeitos colaterais.
Modelo de programa
Uma função durável é um grafo criado a partir de etapas, operadores e funções internas. Strings SQL simples são encapsuladas automaticamente, portanto, você não precisa chamar df.sql() explicitamente.
Operators
| Operador | Meaning | Exemplo |
|---|---|---|
~> |
Sequência – executar à esquerda, depois à direita | 'SELECT 1' ~> 'SELECT 2' |
& |
Junção – executar em paralelo, aguardar todos | 'SELECT 1' & 'SELECT 2' |
| |
Corrida – executar em paralelo, primeiro vence | fast_query | df.sleep(30) |
?>
!>
|
If/else – ramificação com base em uma condição booliana | cond ?> then_branch !> else_branch |
@> |
Loop – repetição para sempre (operador prefixado) | @> body |
|=> |
Nome - registrar o resultado de uma etapa | 'SELECT id FROM users LIMIT 1' |=> 'user_id' |
Recursos integrados úteis
| Função | Purpose |
|---|---|
df.sleep(seconds) |
Pausar por N segundos. Durável entre reinicializações. |
df.wait_for_schedule(cron) |
Aguarde até a próxima vez na qual a expressão cron seja atendida. |
df.wait_for_signal(name, timeout) |
Bloqueie até que um df.signal() externo chegue. |
df.http(url, method, body, headers, timeout) |
Faça uma chamada HTTP como uma atividade durável, com nova tentativa em caso de falha transitória. |
df.if(cond, then, else) |
Ramificação condicional. |
df.loop(body, cond) |
Repita enquanto uma condição SQL é verdadeira. |
df.join(a, b) / df.race(a, b) |
Execução paralela e de corrida. |
df.join3(a, b, c) |
Para execução paralela em três vias. |
df.start(body, label, database) |
Envie uma função durável e retorne o identificador da instância. |
df.cancel(id, reason) |
Cancelar uma instância em execução. |
df.status(id) / df.result(id) |
Inspecione o resultado. |
df.explain(input) |
Renderize o grafo de funções para visualização. |
Leia mais sobre todos os recursos do pg_durable.
Variáveis
|=> captura o resultado de uma etapa com um nome; etapas posteriores fazem referência a ele como $name.
SELECT df.start(
'SELECT 100 AS amount' |=> 'total'
~> 'SELECT $total * 2 AS doubled'
);
Exemplos de uso
ETL multi-etapa com novas tentativas
Um ETL diário que limpa, carrega, indexa 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 o banco de dados for reiniciado entre o DELETE e o INSERT, o trabalho será reiniciado no INSERT – ele não reexecuta o DELETE.
Tarefa agendada (cron)
Execute uma tarefa de manutenção todos os dias da semana às 9h:
SELECT df.start(
@> (
df.wait_for_schedule('0 9 * * 1-5')
~> 'CALL refresh_materialized_views()'
),
'weekday-refresh'
);
Se você quiser interromper esse trabalho, poderá executar a cancel função.
SELECT df.cancel('a1b2c3d4', 'stop test cron job');
Fluxo de trabalho de aprovação com tempo limite
Aguarde até 24 horas por um sinal de aprovação externo e, em seguida, confirme ou rejeite:
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"}');
Chamada HTTP durável
df.http() faz chamadas externas como atividades duráveis – Respostas 5xx, erros de rede e tempos limite são repetidos 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'
);
Leia mais sobre a segurança HTTP permitida no pg_durable.
Observar e operar
Tudo pode ser consultado do SQL. Não há nenhuma interface do usuário ou serviço separado para aprender.
-- 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();
Para verificar se o trabalhador está vivo:
SELECT epoch_id, last_seen_at, now() - last_seen_at AS time_since_last_heartbeat
FROM df._worker_epoch;
Menos de time_since_last_heartbeat 15 segundos significa que o trabalhador está saudável. Qualquer coisa maior ou nenhuma linha indica que o trabalho permanece inativo ou não foi inicializado.
Monitorar fluxos de trabalho no Visual Studio Code
A extensão do PostgreSQL para o Visual Studio Code inclui uma aba Fluxos de trabalho na exibição Pipelines & Workflows, em que você pode inspecionar pg_durable instâncias de fluxo de trabalho e monitorar o estado da execução no editor.
Abrir o painel Fluxos de Trabalho
- Em Visual Studio Code, abra a extensão PostgreSQL.
- Em Pesquisador de Objetos, clique com o botão direito do mouse no banco de dados.
- Selecione Pipelines &Fluxos de Trabalho.
- Selecione a guia Fluxos de Trabalho .
O painel esquerdo lista PG Durable Runs e o painel central mostra detalhes da instância de fluxo de trabalho selecionada.
Inspecionar execuções de fluxo de trabalho
Ao selecionar uma execução de fluxo de trabalho, examine o resumo para validar:
-
Status:
completed,runningoufailed. - ID de execução: identificador exclusivo para a instância.
- Tempo e duração iniciados: acompanhe o progresso e o desempenho da execução.
- Painel de detalhes: metadados de execução adicionais.
Use as abas disponíveis para explorar mais a fundo:
- Gráfico: Exibição visual da execução passo a passo, mostrando a estrutura do fluxo de trabalho e o fluxo das etapas.
- Temporização: visualização focada na duração para análise do desempenho e identificação de gargalos.
- Resultados: Saídas e detalhes orientados a resultados da execução do fluxo de trabalho.
Para fluxos de trabalho relacionados a pipelines de IA, a ação Ver definição do pipeline (quando disponível) permite acessar, a partir de uma execução do fluxo de trabalho, a definição do pipeline, o que é útil para comparar o comportamento entre execuções ou investigar regressões.
Identidade e isolamento
As funções duráveis são executadas com os privilégios do usuário que as enviou, não com os privilégios do trabalhador.
pg_durable captura tanto session_user quanto current_user no momento do envio, de modo que as funções enviadas em um contexto SET ROLE sejam executadas com esse papel efetivo.
Isso significa que:
- Os usuários só veem e modificam dados que já têm permissões para acessar.
- Os não superusuários não podem escalonar privilégios enviando uma função durável.
- As cargas de trabalho multilocatário permanecem isoladas desde que a função e o modelo de concessão estejam corretos.
Interação com réplicas, backup e PITR
- Backup e PITR. O grafo de funções (
dfesquema) e o estado de execução (duroxideesquema) são armazenados em tabelas regulares e estão incluídos nos backups do HorizonDB. Uma restauração pontual restaura ambos. - Réplicas de leitura. O trabalho em segundo plano só é executado no primário. As réplicas de leitura podem consultar as
df.*exibições de monitoramento, mas não executam funções. - Failover. Depois de um failover, o trabalho no novo primário continuará de onde o primário anterior parou. As instâncias em execução são retomadas no ponto de verificação mais recente.
Comparado a orquestradores externos
| Aspecto | Orquestrador externo | pg_durable |
|---|---|---|
| Implantação | Serviço separado, identidade separada, repositório de estado separado | Um banco de dados |
| Durabilidade do estado | Camada de armazenamento do Orchestrator | Os mesmos backups, HA e PITR dos dados |
| Identity | Os trabalhos são executados em uma identidade do serviço | As funções são executadas como o usuário que as enviou |
| Modos de falha | Rede entre o orquestrador e o banco de dados | Nenhum - mesmo processo |
| Melhor para | Orquestração entre sistemas que envolve muitos serviços | Cargas de trabalho em que a maior parte do trabalho está no Postgres ou próximo |
pg_durable não está tentando substituir orquestradores externos em pipelines entre sistemas diferentes. É a escolha certa quando a maior parte do trabalho é o trabalho do banco de dados - inserções, transformações, chamadas de IA, manutenção agendada - e adicionar outro serviço é mais custo do que benefício.
Limitações durante a visualização
-
df.http()faz novas tentativas em erros 5xx e de rede. As respostas 4xx são retornadas ao fluxo de trabalho para você tratá-las; não são tentadas novamente automaticamente. - O trabalho em segundo plano atende a um único banco de dados por instância. O fan-out de vários bancos de dados é compatível por meio de
df.start(..., database => 'other_db')em uma função em execução no banco de dados do trabalho. - As definições de função e o estado de execução não são portáveis entre versões principais de
pg_durabledurante a fase de visualização. Escorra ou cancele as instâncias em execução antes de atualizar.
Conteúdo relacionado
- Implementar pipelines de IA duráveis no Azure HorizonDB (versão prévia)
- Funções de IA na extensão azure_ai para o Azure HorizonDB (Versão Prévia)
- Gerar inserções de vetor usando a função de IA create_embeddings() (versão prévia)
- Permitir extensões no Azure HorizonDB (versão prévia)
- Duroxide no GitHub