Funções duráveis com pg_durable para Azure HorizonDB (Pré-visualização)

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 df e duroxide. 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 depois CREATE EXTENSION e 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:

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

  1. No Visual Studio Code, abra a extensão PostgreSQL.
  2. No Object Explorer, clique com o botão direito na sua base de dados.
  3. Selecione Pipelines & Fluxos de Trabalho.
  4. Selecione o separador Workflows.

O painel esquerdo lista PG Durable Runs, e o painel central mostra detalhes para a instância de workflow selecionada.

Captura de ecrã do separador Fluxos de trabalho na extensão PostgreSQL para Visual Studio Code, a mostrar PG Durable Runs e detalhes do fluxo de trabalho.

Inspecionar execuções de fluxo de trabalho

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

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