Ir al contenido principal
Socios

Construye agentes duraderos con Temporal y Lakebase

por Sam Ingbar

  • Preserve el progreso del agente con Temporal: recupere el trabajo registrado tras fallos de los workers, reintente las operaciones fallidas y espere de forma duradera la revisión humana.
  • Sirva el estado de la aplicación en tiempo real con Lakebase Postgres: haga que las pruebas, recomendaciones, decisiones de revisión y métricas operativas se puedan consultar durante cada ejecución.
  • Conecte la ejecución con datos gobernados: lea las políticas de Unity Catalog a través de tablas sincronizadas y, con Lakebase Change Data Feed habilitado, publique los cambios operativos de vuelta en las tablas de historial de Delta.

Un agente de suscripción de préstamos personales recopila información, aplica políticas y puede esperar días a que un revisor responda. Durante ese tiempo, los workers pueden reiniciarse y las llamadas a herramientas pueden fallar. La aplicación debe conservar el trabajo completado, reanudar la ejecución y mantener la información disponible para el revisor.

Esta implementación de referencia utiliza Temporal para una ejecución duradera y Lakebase Postgres para un estado operativo consultable. Una tabla sincronizada hace que la política de suscripción de Unity Catalog esté disponible en Lakebase. Las actividades de Temporal escriben información, decisiones y métricas en Lakebase; una vez habilitado, el Change Data Feed de Lakebase puede publicar esos cambios en las tablas de historial de Delta administradas por Unity Catalog. Esta combinación es especialmente útil cuando Databricks ya administra las entradas del agente y el análisis posterior.

El desafío de los agentes en la nube de larga duración

Un agente en la nube puede durar más que la solicitud, el worker, el contenedor o el despliegue que lo inició. Un usuario puede iniciar una sesión, regresar al día siguiente y continuar en otro worker. Los despliegues y los fallos de proceso son habituales, por lo que el progreso del agente debe sobrevivir de forma independiente al proceso que lo ejecuta. La recuperación requiere tanto los resultados de las operaciones completadas como el estado del flujo de control necesario para determinar qué sucede a continuación.

Para este agente de suscripción, esto genera seis requisitos:

  1. Recuperación: Un worker de reemplazo debe reanudar la ejecución desde el último paso completado.
  2. Reintentos: Las llamadas a herramientas y las operaciones de base de datos deben tolerar la ejecución repetida sin duplicar efectos secundarios.
  3. Largas esperas: El agente debe esperar a las personas o a los sistemas externos sin mantener un worker abierto.
  4. Visibilidad operativa: Las aplicaciones y los operadores necesitan conocer el estado actual, la información recopilada, el estado de los reintentos y los detalles de los fallos.
  5. Gobernanza en tiempo de ejecución: Las actualizaciones de las políticas deben estar disponibles sin necesidad de desplegar código, y la aplicación debe definir cuándo las adopta un caso abierto.
  6. Auditoría: El sistema debe conservar la información, la política, la recomendación y la decisión humana asociadas a cada ejecución.

La transcripción de una conversación solo cubre parte de este estado. La recuperación también requiere el historial del flujo de control: qué operaciones se programaron, qué resultados se registraron, qué está esperando el agente y qué comandos ha aceptado.

Temporal simplifica la gestión de sistemas distribuidos. Al compilar con Temporal, un Workflow es el flujo de control duradero para la ejecución de un agente. Una Activity es una llamada a un modelo, herramienta o base de datos cuyo resultado se registra en el historial de eventos del Workflow; las Activities se pueden reintentar. Una Signal es un comando asíncrono enviado a un Workflow en ejecución, como la decisión de un evaluador de riesgos. La implementación de referencia de Temporal Lakebase AgentWorkflow es un agente ejecutable de suscripción de préstamos personales. Llama a varias herramientas, lee políticas gobernadas, genera una recomendación y espera a un evaluador de riesgos.

Lakebase Postgres también ayuda a los desarrolladores a gestionar estos problemas, pero Temporal y Lakebase almacenan estados diferentes para distintos consumidores. El historial de eventos de Temporal impulsa la reproducción. Lakebase almacena la vista orientada a la aplicación: el estado actual de la ejecución, los mensajes, la información recopilada, el estado de la revisión y las métricas. Unity Catalog sigue siendo la fuente de políticas; una tabla sincronizada hace que esa política se pueda consultar en Postgres, y Change Data Feed proporciona la ruta de retorno para el historial operativo. Los sistemas no comparten una transacción. Las escrituras de Lakebase se ejecutan como Temporal Activities bajo una ejecución de al menos una vez (at-least-once). Los identificadores deterministas, las restricciones, las actualizaciones protegidas y los upserts de Postgres garantizan que los intentos repetidos de Activity se dirijan al mismo registro lógico.

Esta arquitectura añade dos sistemas administrados y un contrato de proyección entre ellos. Juntos, mejoran la resiliencia y la escalabilidad del agente al tiempo que mantienen baja la sobrecarga operativa. Temporal más Lakebase es muy útil cuando una sesión de agente debe sobrevivir al reemplazo de un worker, aceptar entradas después de largas esperas, exponer el estado relacional a una aplicación y aplicar datos gobernados mientras permanece abierta.

El caso de uso de suscripción

Elegí la suscripción de préstamos porque la misma ejecución debe recopilar información, aplicar políticas, generar una recomendación y esperar a una persona. Un worker puede fallar entre cualquiera de esos pasos. La política puede cambiar sin necesidad de desplegar la aplicación, y la UI necesita la información actual antes de que se cierre el Workflow.

Los solicitantes simulados reemplazan a las agencias de crédito y a los proveedores de ingresos reales, y la secuencia de herramientas es determinista para mayor simplicidad. Cada solicitud contiene un ID de usuario, un ID de solicitante, el importe, el propósito, la elección del modelo y un límite de turnos. FastAPI asigna run_id, inicia LoanUnderwritingWorkflow y utiliza ese mismo ID en la API, la ejecución de Temporal y las filas de Lakebase.

En el primer turno, credit_check devuelve la puntuación, las líneas de crédito, los impagos y la deuda actual. income_verification devuelve información sobre ingresos y empleo. debt_to_income_calc calcula la relación deuda-ingresos. policy_lookup carga la política para el propósito del préstamo y evalúa la información frente a los umbrales de aprobación, remisión y rechazo definitivo.

El solicitante límite de la muestra tiene una puntuación de crédito de 665, 76 000 USD de ingresos anuales verificados, 2400 USD de deuda mensual y un indicador de impago no material. El resultado de la política registra cada regla, umbral, valor real, resultado de aprobado/reprobado, fuente, recomendación y justificación. El modelo puede recomendar, pero no decidir. Un evaluador de riesgos aprueba, deniega o solicita más información. Una solicitud de más información se convierte en otro mensaje de usuario y en otro turno del agente. El caso pone a prueba una caída del worker después de llamadas a herramientas completadas, una escritura confirmada en Lakebase cuya finalización de Activity se ha perdido, una revisión que se deja abierta durante días, una decisión del navegador obsoleta y un cambio de política durante la ejecución.

Arquitectura

image1.jpg
Figura 1. Las rutas de ejecución, estado operativo y gobernanza en la implementación de referencia.

Para implementar el agente de suscripción, React y FastAPI se encargan del trabajo de HTTP y UI: iniciar ejecuciones, renderizar información, enumerar casos y enviar decisiones de revisión. Temporal Cloud almacena el historial de eventos y distribuye las tareas (Tasks). Los workers reproducen el código del Workflow y ejecutan las Activities de modelo, herramienta y Lakebase; la I/O de red y de base de datos permanece fuera del código determinista del Workflow.

Una ejecución comienza cuando FastAPI inicia un Workflow. El worker programa las Activities, Temporal registra sus resultados y el agente finalmente llega a AWAITING_REVIEW. La respuesta del evaluador de riesgos se devuelve a través de una Signal. La aprobación o denegación cierra la ejecución; una solicitud de más información reanuda el bucle del agente.

Lakebase contiene dos esquemas operativos. agent_ops contiene el estado de la ejecución, los mensajes, las llamadas a herramientas, los registros de revisión, los eventos y las métricas que FastAPI puede consultar con SQL. agent_policy contiene la política sincronizada de solo lectura utilizada por policy_lookup. Cada Activity escribe registros indexados por los mismos identificadores deterministas utilizados por el Workflow, de modo que la proyección puede ponerse al día después de un reintento sin que Lakebase forme parte del mecanismo de reproducción de Temporal.

Unity Catalog es la fuente de los umbrales de suscripción. Una tabla sincronizada continua los pone a disposición del agente en ejecución. Los umbrales aplicados, la información recopilada y la decisión humana posterior se escriben en agent_ops. Change Data Feed puede publicar esos cambios en las tablas de historial administradas por Unity Catalog para su auditoría y análisis.

Recuperar el trabajo completado después de un fallo del worker

Temporal conserva el historial de eventos ordenado necesario para reconstruir el estado del Workflow en otro worker. Ese historial incluye la programación y los resultados de las Activities, los temporizadores y las Signals. La reproducción ejecuta el código del Workflow con respecto a esos eventos registrados y reconstruye variables como el turno actual, las decisiones de revisión aceptadas, el uso de tokens y la información recopilada.

Durante la reproducción, se devuelve un resultado de Activity registrado en lugar de volver a ejecutar la Activity. Una verificación de crédito completada sigue estando completada, y una respuesta del modelo registrada sigue siendo la respuesta para esa ejecución. Si una Activity estaba en curso cuando falló el worker y Temporal nunca registró su finalización, Temporal puede programar otro intento. Para un agente, esto conserva las respuestas del modelo que ya están registradas en el historial de eventos. Una llamada al modelo cuya finalización no se registró puede volver a ejecutarse, incluso si el proveedor terminó de procesarla.

Las políticas de reintento (Retry Policies) se asignan con la granularidad de las operaciones individuales y se pueden reutilizar en el código. En el ejemplo, las actividades (Activities) que llaman al modelo permiten hasta cuatro intentos dentro de un tiempo de espera de programación a cierre (schedule-to-close timeout) de tres minutos. Las actividades de llamada a herramientas permiten hasta tres intentos y tienen un tiempo de espera de inicio a cierre (start-to-close timeout) de 60 segundos. Las actividades de Lakebase permiten hasta cinco intentos con un tiempo de espera de inicio a cierre de 15 segundos.

Hacer que los efectos externos sean seguros de repetir

Un riesgo es que una escritura de resultado de herramienta de Lakebase se confirme (commit) antes de que el Worker informe la finalización de la actividad (Activity). Si la conexión se cae en ese intervalo, Temporal no tiene ningún resultado registrado y programa otro intento. Ambos intentos representan la misma escritura lógica.

Cada registro de Lakebase tiene una identidad estable. run_id ancla el esquema operativo. message_id identifica un mensaje, tool_call_id una invocación de herramienta, event_id un hito, review_id una ronda de revisión y decision_id un comando de revisor. Las claves primarias y las restricciones de unicidad (unique constraints) de Postgres aplican esas identidades.

La escritura de inicio de herramienta muestra tanto la identidad estable como la protección de estado terminal (terminal-state guard):

Un reintento se dirige al mismo tool_call_id. El predicado final solo permite que una fila no terminal existente se vuelva a escribir como iniciada (started). Si la fila ya tiene un estado de éxito (succeeded) o fallo (failed), PostgreSQL afecta a cero filas. No genera ningún error.

El llamador debe inspeccionar un resultado de cero filas. LakebaseWriteResult devuelve el recuento de filas afectadas, pero el wrapper de actividad (Activity wrapper) actual no convierte el cero en un fallo. El código de producción debería clasificar el cero como una operación nula (no-op) esperada solo después de confirmar el estado terminal almacenado; de lo contrario, debería generar o registrar un conflicto. La misma regla se aplica a las transiciones protegidas de ejecución y revisión.

Actualizaciones/inserciones (upserts) similares cubren mensajes, resultados de herramientas y eventos (Events). Los ID deterministas hacen que los reintentos converjan en la misma fila lógica, mientras que cada escritura protegida define qué transiciones de estado son legales. La API puede mostrar brevemente un estado anterior mientras se reintenta una escritura. Después de que la actividad (Activity) tiene éxito, la fila aceptada se puede consultar.

Cada herramienta con efectos secundarios necesita un contrato equivalente. Una API de pago puede aceptar una clave de idempotencia, un servicio de correo electrónico un ID de mensaje proporcionado por el llamador y una base de datos una restricción de unicidad. Si el sistema externo no proporciona ningún mecanismo de deduplicación, la actividad (Activity) necesita su propio registro o un proceso de conciliación. Temporal determina cuándo reintentar. La actividad (Activity) determina cómo maneja el sistema externo ese reintento.

Exponer el estado actual y las métricas operativas

El historial de eventos (Event History) proporciona semántica de ejecución y detalles de depuración. La aplicación necesita consultas relacionales indexadas sobre la ejecución actual: listar casos por usuario y estado, cargar una transcripción con su evidencia, buscar revisiones que esperan a una persona y agregar mediciones entre ejecuciones.

Lakebase almacena esa vista de la aplicación en un esquema de Postgres normalizado. agent_runs contiene el estado actual, el ID de flujo de trabajo (Workflow ID), la solicitud, los totales de tokens, las marcas de tiempo y los metadatos de recomendación. agent_messages conserva la transcripción. agent_tool_calls registra los argumentos, el estado, el resultado estructurado, el error y los tiempos. agent_review_decisions conecta la recomendación con un ID de revisión estable, el comando del revisor, la justificación y la hora de la decisión.

El esquema también registra eventos (Events) con nombre y métricas a nivel de flujo de trabajo (Workflow), turno e intento de actividad (Activity-attempt). FastAPI expone los endpoints run-detail, workflow-metrics y retry-metrics respaldados por estas tablas. La interfaz de usuario (UI) puede mostrar una ejecución recopilando evidencia, otra esperando revisión y una tercera reintentando una herramienta fallida. Los operadores pueden consultar las mismas filas con SQL.

La evidencia está disponible antes de que se complete el flujo de trabajo (Workflow). Después de que policy_lookup finaliza, su resultado estructurado se almacena con la llamada a la herramienta. Cuando la ejecución llega a AWAITING_REVIEW, el evaluador de riesgos (underwriter) puede ver la puntuación de crédito, la relación deuda-ingresos (DTI), los umbrales, los resultados de las reglas, la justificación y la fuente de la política que produjo la recomendación.

Mantener la revisión humana duradera y rechazar comandos obsoletos

Cuando el modelo devuelve una recomendación, el flujo de trabajo (Workflow) deriva review_id a partir de run_id y el turno actual. Escribe la revisión pendiente en Lakebase, registra un evento agent.review_pending, establece la proyección en AWAITING_REVIEW y llama a workflow.wait_condition. Temporal conserva el flujo de trabajo (Workflow) abierto sin mantener ocupado un proceso de Worker.

La API envía la acción del evaluador de riesgos como una señal (Signal). Antes de enviarla, la API verifica que Lakebase muestre la ejecución en espera de revisión y que el review_id enviado coincida con la ronda actual. Si alguna de las verificaciones falla, la API devuelve un conflicto. El flujo de trabajo (Workflow) valida de forma independiente el comando contra su propio estado e ignora las decisiones obsoletas o duplicadas, protegiendo la ejecución incluso cuando la proyección de Lakebase se retrasa.

Después de aceptar la señal (Signal), el flujo de trabajo (Workflow) persiste la decisión a través de una actividad (Activity) de Lakebase idempotente. La aprobación o denegación completa la ejecución. Una solicitud de más información cambia la proyección de nuevo a RUNNING, añade la justificación del revisor como un mensaje de usuario e inicia el siguiente turno. Debido a que el turno cambió, la siguiente recomendación recibe un nuevo review_id.

La respuesta 202 de la API confirma que Temporal recibió la señal (Signal). La aceptación comercial ocurre de forma asíncrona en el flujo de trabajo (Workflow), por lo que un comando puede superar la verificación previa de la API y aun así ser ignorado si el estado de la revisión ha cambiado. El cliente actualiza la proyección de Lakebase para observar el estado resultante.

Servir políticas gobernadas sin volver a implementar Workers

Los umbrales de evaluación de riesgos cambian independientemente del código del Worker. La tabla de origen en Unity Catalog contiene valores específicos para cada propósito, como la puntuación de crédito mínima, la relación deuda-ingresos (DTI) de aprobación automática, los umbrales de rechazo directo y el nombre de la política.

El script de configuración crea una tabla sincronizada continua de Lakebase llamada agent_policy.underwriting_policy_limits. policy_lookup que consulta esta copia de Postgres de solo lectura por propósito de préstamo normalizado. Los propietarios de las políticas actualizan el origen de Unity Catalog; la canalización de sincronización propaga el cambio y una ejecución posterior lo lee sin necesidad de implementar un Worker o una API.

El resultado de la política contiene los umbrales aplicados, el valor real de cada regla, el resultado de aprobado/reprobado y el origen. La demostración puede recurrir a una política de respaldo (fixture policy) cuando Lakebase está deshabilitado o la fila no está disponible, y registra esa ruta como fixture_fallback. Un flujo de trabajo (Workflow) regulado puede, en su lugar, fallar cerrado (fail closed). La aplicación tiene que tomar esa decisión de respaldo de forma explícita.

Devolver los cambios operativos a Unity Catalog

El repositorio prepara cada tabla agent_ops para el Change Data Feed de Lakebase estableciendo REPLICA IDENTITY FULL. Un administrador todavía tiene que habilitar la función para el esquema. Luego, Lakebase captura las inserciones, actualizaciones y eliminaciones del registro de escritura anticipada (write-ahead log) de Postgres y las escribe en lotes en tablas de historial de Delta administradas por Unity Catalog con el patrón lb_<table>_history.

El Change Data Feed se encuentra actualmente en vista previa pública (Public Preview) y transmite los cambios aproximadamente cada 15 segundos. Ese intervalo es adecuado para auditorías y análisis, mientras que la interfaz de usuario (UI) consulta directamente a Lakebase para obtener el estado operativo actual.

Las tablas de historial pueden reconstruir el origen de la política de una ejecución, la evidencia de la herramienta, los intentos de actividad (Activity), la espera de revisión, la recomendación y la decisión humana. El repositorio configura el esquema de origen para esta ruta, pero no incluye una ejecución de Change Data Feed de extremo a extremo observada. Habilitar el feed y verificar las tablas de destino siguen siendo pasos de implementación.

Operar el sistema

La implementación separa React/FastAPI del Worker de Temporal. Las réplicas de la API se escalan con la carga de solicitudes; los Workers se escalan con el flujo de trabajo (Workflow) y el backlog de tareas de la actividad (Activity Task backlog) y la concurrencia configurada. El escalado automático de Lakebase ajusta el cómputo de la base de datos dentro de los límites del proyecto.

El precio de Temporal Cloud se basa en las acciones (Actions) más el almacenamiento del historial de eventos (Event History) activo y retenido, por lo que la frecuencia de reintentos y los historiales abiertos durante mucho tiempo también afectan al costo. El equipo aún define las réplicas de Kubernetes, las colas de tareas (Task Queues), los límites del pool de conexiones (connection-pool) y los límites de la base de datos para su carga de trabajo. Alternativamente, puede configurar su propio servicio de Temporal de código abierto utilizando la última versión de código abierto.

El cliente de Lakebase utiliza autenticación de máquina a máquina OAuth. Los tokens de OAuth de Databricks y las credenciales de base de datos generadas caducan, por lo que el cliente actualiza su pool de conexiones SQLAlchemy antes de que caduque la credencial de la base de datos de una hora. Las conexiones utilizan TLS. Sin rotación, un Worker de larga ejecución encontraría fallos en la base de datos en un horario predecible.

Los operadores utilizan Temporal para inspeccionar el historial de flujos de trabajo (Workflow) y actividades (Activity), Lakebase para consultar el estado y las métricas de la aplicación, y Kubernetes para verificar el estado del proceso y de la implementación. Un operador puede entonces distinguir una espera de revisión deliberada de un reintento de actividad (Activity), un fallo de acceso a la base de datos o una herramienta fallida.

Evidencia y límites

El conjunto de pruebas contiene 21 pruebas superadas para la secuenciación de Workflows, el comportamiento de revisión, la construcción de conexiones OAuth, la persistencia idempotente, los contratos de métricas, el inicio de Workflows de API y la configuración de los Workers. El script de recuperación ante fallos añade un ejercicio de fallo de proceso con el proveedor determinista.

Los datos del solicitante y del proveedor son fixtures. El repositorio no valida los modelos de préstamo, el cumplimiento normativo, los controles de seguridad de producción, la disponibilidad regional ni el rendimiento a escala. El ejercicio de fallo local se ejecutó con Lakebase deshabilitado, por lo que aísla la recuperación de Temporal. Change Data Feed aún requiere habilitación y verificación en el entorno de Databricks de destino.

Próximos pasos

¿Quieres saber más? Prueba a ejecutar la demostración por ti mismo y mira la ejecución duradera en acción. Ejecuta la implementación de referencia, detén un worker a mitad de la ejecución y observa cómo se recupera el agente. Conecta Lakebase para explorar las pruebas, las políticas y el estado de revisión humana detrás de cada decisión.

(Esta entrada del blog ha sido traducida utilizando herramientas basadas en inteligencia artificial) Publicación original

Recibe las últimas publicaciones en tu bandeja de entrada

Suscríbete a nuestro blog y recibe las últimas publicaciones directamente en tu bandeja de entrada.