Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
pg_durable é o motor de execução durável dentro Azure HorizonDB. Permite-te definir fluxos de trabalho SQL de longa duração e vários passos (embedding pipelines, jobs ETL, chamadas de IA, jobs agendados, fluxos de aprovação) e executá-los com as mesmas garantias de fiabilidade que esperarias de um orquestrador dedicado como Durable Functions, sem sair do Postgres.
pg_durable é também a camada de execução subjacente a pipelines de IA duradouros. Se estiver a utilizar pipelines de IA, pg_durable é o que lhes permite sobreviver a falhas, repetir em caso de falha e retomar a partir da última etapa concluída.
Note
pg_durable está em versão preliminar.
O que significa "durável"
Uma função duradoura em pg_durable é mantida em disco em cada passo do caminho. Isso dá-lhe um conjunto específico de garantias que não obtém de um simples bloco BEGIN ... COMMIT nem de uma tarefa cron:
- Sobrevive a falhas e reinicios da base de dados. Os passos concluídos não são reexecutados quando o servidor volta a funcionar. Os passos em andamento são retomados a partir do último ponto de verificação. Os passos pendentes são executados quando o worker volta a estar online.
- Resiste a longos períodos de espera. Um fluxo de trabalho pode dormir durante horas, esperar por um cron schedule, ou bloquear um sinal externo, e ainda assim continuar de onde ficou.
- Resiste a falhas. As etapas que falharam podem ser repetidas automaticamente sem voltar a executar a função inteira.
- Captura a identidade. Uma função executa-se com os privilégios do utilizador que a iniciou, não com os privilégios do trabalhador em segundo plano. As cargas de trabalho multitenant mantêm-se isoladas.
- Mantém-se observável a partir do SQL. Pode inspecionar o estado, o histórico, a contagem de execuções e os resultados na mesma interface que usa para tudo o resto no HorizonDB: uma instrução
SELECT.
O que a durabilidade não faz automaticamente: não torna, por si só, seguras para repetir as operações externas não idempotentes. Se uma etapa chamar uma API externa que implica custos, conceba-a de modo a ser idempotente (por exemplo, passando uma chave de idempotência).
Quando usar pg_durable
Utilize pg_durable quando tiver de trabalhar nisso:
- Leva tempo suficiente para falhar a meio do processo (geração de embeddings em milhões de linhas, um processo ETL com várias etapas, um preenchimento retroativo).
- Tem de ser repetido em caso de falha, sem refazer as partes que já foram concluídas com êxito.
- Tem de ser executado de acordo com uma programação (de hora a hora, todos os dias úteis às 9h).
- Precisa de esperar por um evento externo (uma aprovação, um webhook, um sinal de outro sistema).
- Coordena múltiplos passos com ramificação, junção ou corrida.
- Está atualmente implementado como um orquestrador externo + uma base de dados Postgres, onde a maior parte do trabalho é a parte da base de dados.
Se a sua carga de trabalho for uma única instrução transacional curta, não precisa de pg_durable. Use um INSERT / UPDATE normal.
Como funciona
Uma função durável é um gráfico de passos que constróis com um DSL SQL e submetes com df.start(). O grafo é guardado e, em seguida, um processo em segundo plano executa-o.
Duas ideias-chave:
-
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 aplicam-se automaticamente ao estado do seu fluxo de trabalho. Não há um estado orquestrador separado para gerir. - O trabalhador de fundo é iniciado por
shared_preload_libraries. Deteta a extensão depoisCREATE EXTENSIONe começa a executar funções. Se a base de dados reiniciar, o trabalhador volta a ligar-se às instâncias em execução e retoma-as.
Note
O motor de execução dentro do pg_durable é construído sobre Duroxide, o runtime de execução durável open-source da Microsoft para Rust (inspirado no Durable Task Framework e no Temporal). O nome do esquema duroxide reflete isso: é nesse esquema que o Duroxide persiste o histórico de orquestração, os IDs de correlação e o estado de reexecução. As garantias de repetição determinística, ID de evento correlacionado e temporizador durável que se obtêm com pg_durable provêm diretamente do Duroxide.
Ativar pg_durable
Para ativar pg_durable no Azure HorizonDB, configure primeiro um grupo de parâmetros e depois crie a extensão em cada base de dados.
Utilize estes artigos de configuração:
- Cria um grupo de parâmetros para o teu servidor.
- Definir
shared_preload_librariespara incluirpg_durable. - Definir
azure.extensionspara incluirpg_durable. - Aplica o grupo de parâmetros ao servidor.
- Ligue-se a cada base de dados alvo e execute:
Cria a extensão em cada base de dados onde a queres usar:
CREATE EXTENSION IF NOT EXISTS pg_durable;
CREATE EXTENSION prevê o df esquema (grafos de funções e vistas de monitorização) e o duroxide esquema (estado de execução). O trabalhador em segundo plano deteta a extensão em poucos segundos e está pronto para executar funções.
A tua primeira função duradoura
-- 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');
Mesmo uma função de um único passo é persistente: se a base de dados reiniciar depois de df.start() e antes de o processo de trabalho a processar, a função continua a ser executada.
Note
df.start() submete um fluxo de trabalho de forma assíncrona e retorna imediatamente. Para fluxos de trabalho em vários passos, use df.list_instances(), df.instance_info(), df.status(), ou df.result() para confirmar a conclusão antes de validar os efeitos secundários.
Modelo de programa
Uma função durável é um grafo construído a partir de passos, operadores e funções incorporadas. As cadeias SQL simples são automaticamente encapsuladas, pelo que não é necessário invocar df.sql() explicitamente.
Operators
| Operador | Meaning | Example |
|---|---|---|
~> |
Sequência - corre para a esquerda, depois para a direita | 'SELECT 1' ~> 'SELECT 2' |
& |
Junta-te - corre em paralelo, espera por todos | 'SELECT 1' & 'SELECT 2' |
| |
Corrida - corrida em paralelo, primeiras vitórias | fast_query | df.sleep(30) |
?>
!>
|
If / else - ramificação com base numa condição booleana | cond ?> then_branch !> else_branch |
@> |
Loop - repetir para sempre (operador prefixo) | @> body |
|=> |
Nome - registar o resultado de uma etapa | 'SELECT id FROM users LIMIT 1' |=> 'user_id' |
Funcionalidades integradas úteis
| Função | Purpose |
|---|---|
df.sleep(seconds) |
Pausa por N segundos. Durável em reinícios. |
df.wait_for_schedule(cron) |
Espere até à próxima ocorrência em que uma expressão cron corresponda. |
df.wait_for_signal(name, timeout) |
Bloqueia até que chegue um df.signal() externo. |
df.http(url, method, body, headers, timeout) |
Faça uma chamada HTTP enquanto atividade duradoura, com repetição em caso de falhas transitórias. |
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 racial. |
df.join3(a, b, c) |
Para execução paralela de três vias. |
df.start(body, label, database) |
Submeta uma função durável e retorne o respetivo ID de 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 gráfico de funções para visualização. |
Leia mais sobre todas as funcionalidades do pg_durable.
Variáveis
|=> captura o resultado de um passo com um nome; Passos posteriores referem-no como $name.
SELECT df.start(
'SELECT 100 AS amount' |=> 'total'
~> 'SELECT $total * 2 AS doubled'
);
Exemplos de utilização
ETL em várias etapas com novas tentativas
Um ETL diário que limpa, carrega, indexa e regista:
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 a base de dados reiniciar entre a DELETE e a INSERT, o worker retoma a partir da INSERT — não executa novamente a DELETE.
Tarefa agendada (cron)
Execute uma tarefa de manutenção todos os dias úteis às 9h:
SELECT df.start(
@> (
df.wait_for_schedule('0 9 * * 1-5')
~> 'CALL refresh_materialized_views()'
),
'weekday-refresh'
);
Se quiseres parar este trabalho, podes executar a cancel função.
SELECT df.cancel('a1b2c3d4', 'stop test cron job');
Fluxo de aprovação com tempo limite
Aguarde até 24 horas por um sinal externo de aprovação 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"}');
Pedido HTTP duradouro
df.http() faz chamadas externas como atividades persistentes – as respostas 5xx, os erros de rede e os timeouts 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 é consultável a partir de SQL. Não há uma interface 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 time_since_last_heartbeat de 15 segundos significa que o trabalhador está saudável. Qualquer valor superior, ou a ausência total de linhas, significa que o processo está indisponível ou ainda não foi inicializado.
Monitorizar fluxos de trabalho no Visual Studio Code
A extensão PostgreSQL para o Visual Studio Code inclui um separador Workflows na vista Pipelines & Workflows, onde pode inspecionar instâncias de pg_durable fluxo de trabalho e monitorizar o estado de execução a partir do editor.
Abrir o painel de Fluxos de Trabalho
- No Visual Studio Code, abra a extensão PostgreSQL.
- No Object Explorer, clique com o botão direito na sua base de dados.
- Selecione Pipelines & Fluxos de Trabalho.
- Selecione o separador Workflows.
O painel esquerdo lista PG Durable Runs, e o painel central mostra detalhes para a instância de workflow selecionada.
Inspecionar execuções de fluxo de trabalho
Ao selecionar uma execução de fluxo de trabalho, reveja o resumo para validar:
-
Estado:
completed,running, oufailed. - ID da execução: Identificador único para a instância.
- Hora e duração de início: Acompanhar o progresso e o desempenho da execução.
- Painel de detalhes: Metadados adicionais de execução.
Utilize os separadores disponíveis para explorar em maior profundidade:
- Gráfico: Vista visual de execução passo a passo que mostra a estrutura do fluxo de trabalho e o fluxo de passos.
- Temporização: Vista focada na duração para análise de desempenho e identificação de gargalos.
- Resultados: Resultados e detalhes orientados a resultados da execução do fluxo de trabalho.
Para fluxos de trabalho relacionados com pipelines de IA, uma ação Ver definição do pipeline (quando disponível) permite-lhe aceder a partir de uma execução de fluxo de trabalho à respetiva 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 executam-se com os privilégios do utilizador que as submeteu, e não com os privilégios do trabalhador.
pg_durable captura tanto session_user como current_user no momento da submissão, pelo que as funções submetidas num contexto SET ROLE são executadas com essa função efetiva.
Isto significa:
- Os utilizadores só veem e modificam dados a que já têm permissões de acesso.
- Os não-superutilizadores não podem escalar privilégios submetendo uma função duradoura.
- As cargas de trabalho multitenant mantêm-se isoladas enquanto o seu papel e modelo de subsídio estiverem corretos.
Interação com réplicas, backup e PITR
- Backup e PITR. O gráfico de funções (
dfesquema) e o estado de execução (duroxideesquema) são armazenados em tabelas regulares e estão incluídos nas cópias de segurança do HorizonDB. Uma recuperação para um ponto anterior no tempo restaura ambos. - Leia réplicas. O trabalhador em segundo plano só corre na primária. As réplicas de leitura podem consultar as vistas
df.*de monitorização, mas não executam funções. - Failover. Após um failover, o trabalhador no novo primário retoma de onde o primário antigo ficou. As instâncias em execução são retomadas a partir do último ponto de controlo.
Comparado com orquestradores externos
| Aspect | Orquestrador externo | pg_durable |
|---|---|---|
| Deployment | Serviço separado, identidade separada, repositório de estado separado | Uma base de dados |
| Durabilidade do estado | Camada de armazenamento do orquestrador | Mesmos backups, HA e PITR que os teus dados |
| Identity | Os trabalhadores operam sob uma identidade de serviço | As funções são executadas como o utilizador que efetuou a submissão |
| Modos de falha | Rede entre orquestrador e base de dados | Nenhum - mesmo processo |
| Melhor para | Orquestração entre sistemas que abrange muitos serviços | Cargas de trabalho onde a maior parte do trabalho é na Postgres ou perto dela |
pg_durable não está a tentar substituir orquestradores externos para pipelines entre sistemas. É a escolha certa quando a maior parte do trabalho é de bases de dados – embeddings, transformações, chamadas de IA, manutenção programada – e adicionar outro serviço é mais um custo do que um benefício.
Limitações durante a pré-visualização
-
df.http()Tenta novamente o 5xx e erros de rede. As respostas 4xx são devolvidas ao fluxo de trabalho para serem tratadas; não é feita uma nova tentativa automaticamente. - O trabalhador em segundo plano processa uma única base de dados por instância. O fan-out multi-base de dados é suportado através
df.start(..., database => 'other_db')de uma função a correr na base de dados do trabalhador. - As definições de funções e o estado de execução não são portáteis entre versões principais de
pg_durabledurante a pré-visualização. Esvazie ou cancele as instâncias em execução antes de atualizar.
Conteúdo relacionado
- Implementar pipelines de IA duráveis no Azure HorizonDB (Pré-visualização)
- AI funciona na extensão azure_ai para Azure HorizonDB (Preview)
- Gerar embeddings vetoriais usando a função de IA create_embeddings() (Pré-visualização)
- Permitir extensões no Azure HorizonDB (Pré-visualização)
- Duroxide no GitHub