Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
pg_durable es el motor de ejecución duradero dentro de Azure HorizonDB. Permite definir flujos de trabajo sql de ejecución prolongada y de varios pasos (insertar canalizaciones, trabajos ETL, llamadas de IA, trabajos programados, flujos de aprobación) y ejecutarlos con las mismas garantías de confiabilidad que esperaría de un orquestador dedicado, como Durable Functions, sin salir de Postgres.
pg_durable también es la capa de ejecución subyacente a las canalizaciones de IA duraderas. Si usa canalizaciones de IA, pg_durable es lo que hace que superen los fallos, reintenten tras un error y se reanuden a partir del último paso completado.
Note
pg_durable está en versión preliminar.
Qué significa "durable"
Una función duradera en pg_durable se guarda en disco en cada paso del proceso. Eso le da un conjunto concreto de garantías que no se obtienen con un bloque simple BEGIN ... COMMIT ni con una tarea cron:
- Sobrevive a fallos y reinicios de la base de datos. Los pasos ya completados no se vuelven a ejecutar cuando el servidor vuelve a estar en funcionamiento. Los pasos en curso se reanudan desde el último punto de control. Los pasos pendientes se ejecutan cuando el trabajador vuelve a conectarse.
- Resiste largas esperas. Un flujo de trabajo puede quedar en espera durante horas, esperar una tarea programada con cron o quedar bloqueado a la espera de una señal externa, y aun así reanudar la ejecución donde la dejó.
- Resiste fallos. Los pasos con errores se pueden reintentar automáticamente sin volver a ejecutar toda la función.
- Captura la identidad. Una función se ejecuta con los privilegios del usuario que lo inició, no con los privilegios del trabajo en segundo plano. Las cargas de trabajo multicliente se mantienen aisladas.
- Permanece observable desde SQL. Puede inspeccionar el estado, el historial, el recuento de ejecuciones y las salidas a través de la misma interfaz que usa para todo lo demás en HorizonDB: una
SELECTinstrucción .
Lo que la durabilidad no hace automáticamente: no hace que las operaciones externas no idempotentes sean seguras de volver a intentar por sí mismas. Si un paso llama a una API externa que implica un coste, diseñe el paso para que sea idempotente (por ejemplo, pasando una clave de idempotencia).
Cuándo usar pg_durable
Use pg_durable cuando tenga que trabajar así:
- Tarda bastante en fallar en la fase intermedia (generación de incrustaciones en millones de filas, un proceso ETL de varios pasos, una recarga de datos).
- En caso de error, es necesario volver a intentarlo sin repetir las partes que ya se han completado con éxito.
- Debe ejecutarse según una programación (cada hora, cada día de la semana a las 9:00).
- Debe esperar un evento externo (una aprobación, un webhook, una señal de otro sistema).
- Coordina varios pasos con bifurcación, unión o competición.
- Actualmente se implementa como orquestador externo + una base de datos de Postgres, donde la mayor parte del trabajo es la parte de la base de datos.
Si la carga de trabajo consiste en una sola instrucción transaccional breve, no necesita pg_durable. Use uno normal INSERT / UPDATE.
Cómo funciona
Una función duradera es un gráfico de pasos que se compilan con un DSL de SQL y se envían con df.start(). El gráfico se almacena y, a continuación, un proceso en segundo plano lo ejecuta.
Dos ideas clave:
-
El grafo de funciones y el estado de ejecución se almacenan en la propia HorizonDB, en los esquemas
dfyduroxide. Las copias de seguridad, la restauración a un momento dado y la alta disponibilidad se aplican automáticamente al estado del flujo de trabajo. No hay ningún estado de orquestador independiente para administrar. - El proceso en segundo plano se inicia mediante
shared_preload_libraries. Detecta la extensión después deCREATE EXTENSIONy comienza a ejecutar funciones. Si la base de datos se reinicia, el trabajo vuelve a asociarse a las instancias en ejecución y las reanuda.
Note
El motor de ejecución dentro de pg_durable se basa en Duroxide, el entorno de ejecución duradera de código abierto de Microsoft para Rust (inspirado en Durable Task Framework y Temporal). El duroxide nombre del esquema refleja esto: es donde Duroxide conserva el historial de orquestación, los identificadores de correlación y el estado de reproducción. Las garantías de reproducción determinista, de ID de evento correlacionado y de temporizador duradero que ofrece pg_durable proceden directamente de Duroxide.
Habilitar pg_durable
Para habilitar pg_durable en Azure HorizonDB, configure primero un grupo de parámetros y, a continuación, cree la extensión en cada base de datos.
Use estos artículos de configuración:
- Cree un grupo de parámetros para el servidor.
- Establezca
shared_preload_librariespara incluirpg_durable. - Establezca
azure.extensionspara incluirpg_durable. - Aplique el grupo de parámetros al servidor.
- Conéctese a cada base de datos de destino y ejecute:
Cree la extensión en cada base de datos donde quiera usarla:
CREATE EXTENSION IF NOT EXISTS pg_durable;
CREATE EXTENSION aprovisiona el df esquema (gráficos de funciones y vistas de supervisión) y el duroxide esquema (estado de ejecución). El trabajo en segundo plano detecta la extensión en unos segundos y está lista para ejecutar funciones.
Su primera función duradera
-- 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');
Incluso una función de un único paso es duradera: si la base de datos se reinicia después de df.start() y antes de que el trabajador la recoja, la función se sigue ejecutando.
Note
df.start() envía un flujo de trabajo de forma asincrónica y devuelve inmediatamente. Para los flujos de trabajo de varios pasos, use df.list_instances(), df.instance_info(), df.status()o df.result() para confirmar la finalización antes de validar los efectos secundarios.
Modelo de programa
Una función duradera es un grafo creado a partir de pasos, operadores y funciones integradas. Las cadenas SQL simples se envuelven automáticamente, por lo que no es necesario llamar a df.sql() explícitamente.
Operadores
| Operador | Meaning | Example |
|---|---|---|
~> |
Secuencia: ejecute a la izquierda y, a continuación, a la derecha | 'SELECT 1' ~> 'SELECT 2' |
& |
Unir - ejecutar en paralelo, esperar a que todos finalicen | 'SELECT 1' & 'SELECT 2' |
| |
Carrera: se disputa en paralelo; gana el primero | fast_query | df.sleep(30) |
?>
!>
|
If /else: bifurcación en una condición booleana | cond ?> then_branch !> else_branch |
@> |
Bucle - repetir indefinidamente (operador prefijo) | @> body |
|=> |
Nombre: captura el resultado de un paso | 'SELECT id FROM users LIMIT 1' |=> 'user_id' |
Funciones integradas útiles
| Function | Purpose |
|---|---|
df.sleep(seconds) |
Pausar durante N segundos. Duradero tras los reinicios. |
df.wait_for_schedule(cron) |
Espere hasta la próxima vez que coincida una expresión cron. |
df.wait_for_signal(name, timeout) |
Bloquee hasta que llegue un df.signal() externo. |
df.http(url, method, body, headers, timeout) |
Realizar una llamada HTTP como actividad duradera, con reintento en caso de error transitorio. |
df.if(cond, then, else) |
Rama condicional. |
df.loop(body, cond) |
Repita la repetición mientras una condición SQL es cierta. |
df.join(a, b) / df.race(a, b) |
Ejecución paralela y de carrera. |
df.join3(a, b, c) |
Para la ejecución paralela de tres vías. |
df.start(body, label, database) |
Envíe una función duradera y devuelva su identificador de instancia. |
df.cancel(id, reason) |
Cancelar una instancia en ejecución. |
df.status(id) / df.result(id) |
Revise el resultado. |
df.explain(input) |
Representar el gráfico de funciones para su visualización. |
Obtenga más información sobre todas las características de pg_durable.
Variables
|=> captura el resultado de un paso con un nombre; los pasos posteriores hacen referencia a él como $name.
SELECT df.start(
'SELECT 100 AS amount' |=> 'total'
~> 'SELECT $total * 2 AS doubled'
);
Ejemplos de uso
ETL de varios pasos con reintentos
Un ETL diario que limpia, carga, indexa y 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'
);
Si la base de datos se reinicia entre el DELETE y el INSERT, el proceso se reanuda desde el INSERT; no vuelve a ejecutar el DELETE.
Tarea programada (cron)
Ejecute una tarea de mantenimiento cada día de la semana a las 9:00:
SELECT df.start(
@> (
df.wait_for_schedule('0 9 * * 1-5')
~> 'CALL refresh_materialized_views()'
),
'weekday-refresh'
);
Si desea detener este trabajo, puede ejecutar la cancel función .
SELECT df.cancel('a1b2c3d4', 'stop test cron job');
Flujo de trabajo de aprobación con tiempo de espera
Espere hasta 24 horas para una señal de aprobación externa y, a continuación, confirme o rechace:
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"}');
Llamada HTTP duradera
df.http() realiza llamadas externas como actividades duraderas: las respuestas 5xx, los errores de red y los tiempos de espera se reintentan automáticamente.
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'
);
Obtenga más información sobre la seguridad HTTP permitida en pg_durable.
Supervisar y manejar
Todo es consultable desde SQL. No hay ninguna interfaz de usuario o servicio independiente 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 comprobar que el trabajador está activo:
SELECT epoch_id, last_seen_at, now() - last_seen_at AS time_since_last_heartbeat
FROM df._worker_epoch;
Un time_since_last_heartbeat valor inferior a 15 segundos significa que el trabajador está en buen estado. Cualquier fila más grande o ninguna fila significa que el trabajo está inactivo o no se ha inicializado.
Supervisión de flujos de trabajo en Visual Studio Code
La extensión de PostgreSQL para Visual Studio Code incluye una pestaña Workflows en la vista Pipelines & Workflows, donde puede inspeccionar instancias de flujo de trabajo pg_durable y supervisar el estado de ejecución desde el editor.
Abra el panel Flujos de trabajo.
- En Visual Studio Code, abra la extensión PostgreSQL.
- En Explorador de objetos, haga clic con el botón derecho en la base de datos.
- Seleccione Canalizaciones y flujos de trabajo.
- Seleccione la pestaña Flujos de trabajo .
En el panel izquierdo se enumeran las ejecuciones duraderas de PG y el panel central muestra los detalles de la instancia de flujo de trabajo seleccionada.
Inspeccionar ejecuciones de flujo de trabajo
Al seleccionar una ejecución de flujo de trabajo, revise el resumen para validar:
-
Estado:
completed,runningofailed. - Identificador de ejecución: identificador único de la instancia.
- Tiempo de inicio y duración: realice un seguimiento del progreso y el rendimiento de la ejecución.
- Panel de detalles: Metadatos de ejecución adicionales.
Use las pestañas disponibles para profundizar más:
- Gráfico: vista de ejecución paso a paso visual que muestra la estructura de flujo de trabajo y el flujo de pasos.
- Temporización: una vista centrada en la duración para analizar el rendimiento e identificar cuellos de botella.
- Resultados: Detalles de salida y centrados en los resultados de la ejecución del flujo de trabajo.
En el caso de los flujos de trabajo relacionados con las canalizaciones de IA, una acción Ver definición de canalización (cuando está disponible) le permite vincular desde una ejecución de flujo de trabajo a su definición de canalización, útil para comparar el comportamiento entre ejecuciones o investigar regresiones.
Identidad y aislamiento
Durable Functions se ejecuta con los privilegios del usuario que los envió, no con los privilegios del trabajo.
pg_durable captura tanto session_user como current_user en el momento del envío, por lo que las funciones enviadas en un contexto SET ROLE se ejecutan con ese rol efectivo.
Esto significa lo siguiente:
- Los usuarios solo ven y modifican los datos a los que ya tienen permisos para acceder.
- Los no superusuarios no pueden escalar privilegios mediante el envío de una función duradera.
- Las cargas de trabajo multiinquilino permanecen aisladas siempre que su modelo de roles y permisos sea correcto.
Interacción con réplicas, copias de seguridad y PITR
- Copia de seguridad y PITR. El gráfico de funciones (
dfesquema) y el estado de ejecución (duroxideesquema) se almacenan en tablas normales y se incluyen en copias de seguridad de HorizonDB. Una restauración a un punto concreto en el tiempo restaura ambos. - Réplicas de lectura. El trabajo en segundo plano solo se ejecuta en el servidor principal. Las réplicas de lectura pueden consultar las vistas de monitorización
df.*, pero no pueden ejecutar funciones. - Conmutación por error. Tras una conmutación por error, el nodo de trabajo del nuevo primario retoma la tarea donde la dejó el antiguo primario. Las instancias en ejecución se reanudan desde su último punto de control.
Comparación con orquestadores externos
| Aspecto | Orquestador externo | pg_durable |
|---|---|---|
| Deployment | Servicio independiente, identidad independiente, almacén de estado independiente | Una base de datos |
| Durabilidad del estado | Capa de almacenamiento del orquestador | Las mismas copias de seguridad, alta disponibilidad y PITR que sus datos |
| Identity | Los procesos de trabajo se ejecutan con una identidad de servicio | Las funciones se ejecutan como usuario de envío |
| Modos de fallo | Red entre orquestador y base de datos | Ninguno: mismo proceso |
| Más adecuado para | Orquestación entre sistemas que abarca muchos servicios | Cargas de trabajo en las que la mayor parte del trabajo está en Postgres o cerca de |
pg_durable no pretende reemplazar a los orquestadores externos para flujos de trabajo entre sistemas. Es la opción correcta cuando la mayor parte del trabajo es el trabajo de base de datos( incrustaciones, transformaciones, llamadas de IA, mantenimiento programado) y la adición de otro servicio es más costo que beneficio.
Limitaciones durante la versión preliminar
-
df.http()vuelve a intentarlo en caso de errores 5xx y de red. Las respuestas 4xx se devuelven al flujo de trabajo para que las gestione; no se reintentan automáticamente. - El trabajo en segundo plano ofrece una base de datos única por instancia. Se admite la distribución ramificada de varias bases de datos mediante
df.start(..., database => 'other_db')desde una función que se ejecuta en la base de datos del trabajo. - Las definiciones de función y el estado de ejecución no son portátiles en las versiones principales de
pg_durabledurante la versión preliminar. Purga o cancelación de instancias en ejecución antes de actualizar.
Contenido relacionado
- Implementar canalizaciones duraderas de IA en Azure HorizonDB (versión preliminar)
- AI funciona en la extensión azure_ai para Azure HorizonDB (versión preliminar)
- Generación de incrustaciones de vectores mediante la función de IA create_embeddings() (versión preliminar)
- Permitir extensiones en Azure HorizonDB (versión preliminar)
- Duroxide en GitHub