Funções duráveis com pg_durable para Azure HorizonDB (versão prévia)

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 SELECT instruçã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 df e duroxide. 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ós CREATE EXTENSION e 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:

  1. Crie um grupo de parâmetros para o servidor.
  2. Defina shared_preload_libraries para incluir pg_durable.
  3. Defina azure.extensions para incluir pg_durable.
  4. Aplique o grupo de parâmetros ao servidor.
  5. 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

  1. Em Visual Studio Code, abra a extensão PostgreSQL.
  2. Em Pesquisador de Objetos, clique com o botão direito do mouse no banco de dados.
  3. Selecione Pipelines &Fluxos de Trabalho.
  4. 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.

Captura de tela da guia Fluxos de Trabalho na extensão PostgreSQL para Visual Studio Code, mostrando as execuções do PG Durable e os detalhes do fluxo de trabalho.

Inspecionar execuções de fluxo de trabalho

Ao selecionar uma execução de fluxo de trabalho, examine o resumo para validar:

  • Status: completed, runningou failed.
  • 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 (df esquema) e o estado de execução (duroxide esquema) 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_durable durante a fase de visualização. Escorra ou cancele as instâncias em execução antes de atualizar.