Ir al contenido principal
Producto

Simplifica la orquestación de agentes de IA con Lakebase Postgres

Cómo CLA creó una solución nativa de Databricks para tareas de larga duración, observabilidad y atribución de costos

por Li Yu, Michelle JanneyCoyle, Jon Cormack, Yarri Bryn, Alec Sorensen y Darshana Nair

  • Cola de tareas lista para escalar en Postgres: Un análisis profundo de los patrones que convierten un par de tablas de Lakebase en una cola duradera, concurrente y resistente a fallos para tareas de agentes de larga duración, sin necesidad de un bróker, caché o planificador.
  • Arquitectura totalmente nativa de Databricks: Un diseño de referencia que integra Lakebase, Databricks Apps, Lakeflow Jobs, MLflow y Unity Catalog Volumes en un pipeline de extremo a extremo para el procesamiento de documentos basado en agentes, sin infraestructura externa que operar.
  • Observabilidad e invariantes en tiempo real: Un análisis profundo sobre el uso de disparadores LISTEN/NOTIFY de Postgres combinados con Server-Sent Events (SSE) para crear un panel de operador de baja latencia que realiza un seguimiento automático de costos y tareas sin sobrecarga adicional.

Introducción

Tradicionalmente, la auditoría es un proceso tedioso que a menudo requiere una revisión detallada de documentos y la extracción de información. Para acelerar este proceso, CLA (CliftonLarsonAllen LLP), una firma líder de servicios profesionales con una creciente presencia global, colaboró con el equipo de Forward Deployed Engineering de Databricks para crear y poner en producción una solución de auditoría basada en agentes. De manera conjunta, desarrollamos una aplicación de procesamiento de documentos que reduce el tiempo de extracción de horas a minutos sin comprometer la calidad. La aplicación está construida completamente sobre Databricks, utilizando Lakebase Postgres, Databricks Apps, Lakeflow Jobs, MLflow y Unity Catalog Volumes. En este blog, nos centramos en un componente clave de ese sistema: la capa de orquestación impulsada por Lakebase.

La capa de orquestación se encarga de coordinar las tareas de larga duración, gestionar los reintentos, atribuir costes y ofrecer visibilidad en tiempo real. Con Lakebase y Databricks Apps, eliminamos la necesidad de contar con una infraestructura independiente para la gestión de colas, la orquestación y la observabilidad.

Lakebase también hace que esta arquitectura sea práctica a escala al separar el almacenamiento del cómputo. A diferencia de las implementaciones tradicionales de Postgres, el cómputo puede escalar según la demanda mientras que el almacenamiento permanece duradero e independiente. Juntas, estas capacidades hacen de Lakebase una base práctica para un patrón de orquestación más sencillo y escalable para cargas de trabajo basadas en agentes de larga duración en Databricks.

Desafíos de orquestación para cargas de trabajo basadas en agentes

El análisis de documentos es una carga de trabajo basada en agentes muy común y de gran volumen. Las empresas de todos los sectores necesitan convertir grandes volúmenes de contratos, facturas, declaraciones financieras y otros documentos en datos estructurados. Ejecutar esto a escala plantea cinco problemas distintos de los sistemas distribuidos:

  • Latencia impredecible por tarea: una factura de dos páginas puede procesarse en segundos, mientras que un contrato de doscientas páginas puede tardar varios minutos, lo que dificulta predecir cuánto tiempo se ejecutará cada tarea individual.
  • Limitación de velocidad (throttling) adaptativa: los endpoints de los modelos de visión y LLM limitan el número de solicitudes y tokens que pueden procesar en un periodo determinado. Enviar cientos de tareas a la vez puede superar esos límites, activar la limitación de velocidad y provocar reintentos repetidos. El orquestador debe limitar proactivamente el trabajo en curso (por recuento de tareas concurrentes, por presupuesto de tokens o ambos) en lugar de depender únicamente de los reintentos reactivos.
  • Priorización de la carga de trabajo: los envíos urgentes no deben retrasarse por lotes voluminosos. La prioridad por tarea garantiza que el trabajo de mayor prioridad (envíos interactivos, solicitudes de nivel premium, reprocesamientos iniciados por el operador) se despache primero.
  • Atribución de costes por tarea: los equipos financieros necesitan atribuir el gasto a tareas, clientes y agentes específicos, desglosado por el uso de tokens de IA y el consumo de cómputo.
  • Visibilidad del progreso en tiempo real: los usuarios que suben cientos de documentos necesitan una vista del progreso en vivo.

Muchas organizaciones combinan varios sistemas especializados para la orquestación y la observabilidad. Cada sistema aporta su propia infraestructura, autenticación, monitoreo y requisitos operativos, junto con el trabajo necesario para integrarlos. Para tareas basadas en agentes independientes y de larga duración, esa sobrecarga es desproporcionada en comparación con la complejidad real de la programación.

La solución nativa de Databricks que desarrollamos cumple con todos los requisitos anteriores utilizando Lakebase como base.

Arquitectura de la solución

Arquitectura de la solución

El stack completo de la aplicación consta exclusivamente de servicios de Databricks:

  • Aplicación web (Databricks Apps). Una interfaz de usuario basada en FastAPI donde los usuarios suben archivos PDF (almacenados en Unity Catalog Volumes) y envían solicitudes de análisis. Las solicitudes se escriben directamente en la tabla de tareas de Lakebase.
  • Lakebase. Una base de datos Postgres con escalado automático que aloja el estado relacional del orquestador en tablas relacionadas: tasks (documentos a analizar, que contienen el estado, la información de arrendamiento y el resultado estructurado), y task_attempts (una fila por intento de ejecución, que captura el ID de ejecución de Databricks Job, el ID de traza de MLflow y los metadatos de coste por intento). Lakebase sirve como la única fuente de verdad para el estado del orquestador.
  • Orquestador (Databricks Apps). Un demonio de trabajo de larga duración y un panel de control del operador. El demonio saca las tareas de la cola de Lakebase, las despacha a la capa de agentes de IA y vuelve a escribir los resultados. El panel de control lee las mismas tablas para mostrar el estado en tiempo real.
  • Agentes de IA (Lakeflow Jobs). Los Lakeflow Jobs ejecutan el trabajo de análisis. Cada Job lee un PDF de Unity Catalog Volumes, lo procesa a través de Procesamiento inteligente de documentos y llamadas de visión/LLM, almacena la salida analizada en Lakebase y notifica al orquestador a través de un webhook. MLflow Tracing captura los detalles de la ejecución, como las llamadas al modelo, el uso de tokens, la latencia y los metadatos de coste.

El flujo de datos entre los componentes es el siguiente. La aplicación web escribe los archivos PDF en Unity Catalog Volumes y las solicitudes de análisis en Lakebase. El orquestador saca las tareas de la cola de Lakebase y despacha los Databricks Jobs a la capa de agentes de IA. Los agentes de IA procesan los documentos, vuelven a escribir los resultados en Lakebase y devuelven la llamada al orquestador con actualizaciones de estado.

Debido a estas capacidades integradas en Databricks, no tuvimos que depender de agentes de mensajes externos (Kafka, Redis), planificadores independientes (Airflow, Temporal) ni capas de almacenamiento en caché dedicadas.

Implementación de la cola de tareas

La cola de tareas está respaldada por dos tablas de Postgres en Lakebase. La tabla tasks contiene una fila por unidad lógica de trabajo, registrando el estado actual de la tarea, la información de arrendamiento, la asignación del agente, la extracción principal y el resultado final. La tabla task_attempts contiene una fila por intento de ejecución, capturando el ID de ejecución de Databricks Job, el ID de traza de MLflow y los metadatos de coste por intento. La relación padre-hijo admite reintentos (una sola tarea puede tener varios intentos) y preserva la observabilidad a nivel de intento para la atribución de costes y la depuración.

Un par de tablas de Postgres por sí solas aún no constituyen una cola de tareas. Cuatro patrones nativos de Postgres las transforman en una cola robusta, concurrente, resistente a fallos y adaptativa a los límites de velocidad, adecuada para cargas de trabajo basadas en agentes de larga duración.

Desencolado concurrente con prioridad

Una consulta de desencolado básica podría seleccionar la siguiente tarea disponible utilizando WHERE status = 'enqueued' and LIMIT batch_size. Aunque esa consulta identifica correctamente una tarea en cola, no es suficiente cuando varios trabajadores están desencolando de forma concurrente. Sin el bloqueo de filas, varios trabajadores podrían seleccionar la misma tarea antes de que se actualice el estado.

Agregar FOR UPDATE SKIP LOCKED hace que la desencolación sea segura para la concurrencia. Cada worker bloquea la fila que selecciona, mientras que otros workers omiten esa fila y continúan con la siguiente tarea disponible. Además, una cláusula ORDER BY priority DESC, created_at garantiza que las tareas de mayor prioridad se seleccionen primero, al tiempo que se preserva el orden FIFO dentro de cada nivel de prioridad.

La instrucción completa, segura para la concurrencia y estable para la prioridad, es:

Recuperación ante fallos mediante bloqueo basado en concesiones

Los workers pueden finalizar a mitad de la tarea debido a la expulsión de la VM, condiciones de falta de memoria (out-of-memory) o eventos de despliegue. Si las tareas finalizadas permanecen marcadas como en proceso, se retendrían indefinidamente. La solución es registrar una concesión (lease) que expira en el momento de la desencolación:

Un limpiador (sweeper) periódico vuelve a encolar cualquier tarea cuyo lease_expires_at haya vencido. Las tareas retenidas por workers finalizados se recuperan automáticamente en cuestión de minutos, sin necesidad de un servicio de coordinación externo.

Regulación consciente de los límites de velocidad

Los endpoints de modelos de visión y LLM suelen aplicar dos cuotas distintas: un límite de solicitudes por segundo y un límite de tokens por minuto (TPM). Una única estrategia de regulación rara vez se adapta a ambos. El orquestador admite tres modos, seleccionados por agente mediante la configuración.

Límite de concurrencia. Un parámetro MAX_CONCURRENT_TASKS limita el número de tareas que el orquestador distribuye de forma concurrente. El límite se aplica en el momento de la desencolación contando las filas PROCESSING actuales en la tabla de tareas:

Si el recuento es igual o superior al límite, no se desencola ninguna tarea nueva. Basar la comprobación en el recuento de filas de la base de datos, en lugar de en el tamaño de la cola del ejecutor local, mantiene la precisión del límite a través de los reinicios de los workers, las recuperaciones de concesiones y los despliegues de múltiples réplicas. Este modo es ideal para endpoints limitados por restricciones de solicitudes por segundo donde el uso de tokens de cada tarea es aproximadamente uniforme.

Presupuesto de tokens. Un parámetro MAX_TPM limita la tasa de tokens proyectada en las tareas en curso. El orquestador estima el recuento de tokens de una tarea y suma la tasa de tokens proyectada en todas las tareas PROCESSING. Solo se desencola una nueva tarea si la suma más los tokens proyectados de la nueva tarea se ajusta al presupuesto.

Límite combinado. Cuando se configuran tanto MAX_CONCURRENT_TASKS y MAX_TPM, el orquestador aplica la restricción que sea más estricta. Este modo gestiona cargas de trabajo que están limitadas por la concurrencia bajo un régimen (muchas tareas cortas y de bajo coste) y limitadas por tokens bajo otro (un único documento muy largo que satura la cuota por minuto).

En los tres modos, la decisión de regulación se toma en el momento de la desencolación dentro de la misma transacción que FOR UPDATE SKIP LOCKED. Una tarea que no cabe dentro de la cuota actual permanece en la cola y se vuelve a evaluar en el siguiente ciclo de desencolación: sin un estado de programación independiente, sin cola de espera en memoria y sin capa de coordinación entre las réplicas de los workers.

Callbacks de webhook idempotentes

Cuando la capa de agentes de IA completa una tarea, envía un callback al orquestador con el resultado. La entrega de callbacks no es de tipo exactamente-once: Databricks puede reintentarlo, las redes pueden interrumpirse y los proxies pueden volver a realizar el envío. El controlador de callbacks está diseñado para ser idempotente: acepta los estados PROCESSING y ENQUEUED, y trata las tareas que ya están en estado terminal como operaciones sin efecto (no-ops). Los payloads idénticos producen resultados idénticos, lo que elimina el riesgo de doble facturación o procesamiento duplicado.

Estos cuatro patrones combinados dan como resultado una cola de tareas que es correcta bajo concurrencia, duradera ante fallos, consciente de los límites de velocidad bajo carga e idempotente ante reintentos. La visibilidad en tiempo real del sistema en ejecución se proporciona mediante un mecanismo independiente que se describe en la siguiente sección.

Panel de control del operador en tiempo real

Cuando hay muchos documentos en curso, los operadores necesitan una visión clara del rendimiento del agente, el estado de las tareas y el coste de la carga de trabajo. No deberían tener que consultar continuamente al orquestador ni depender de una plataforma de métricas independiente. El orquestador integra esta funcionalidad directamente en un único panel de control presentado por la misma Databricks App que ejecuta el demonio del worker.

Capacidades del panel de control

El panel de control orientado al operador presenta un conjunto de métricas operativas que, en conjunto, caracterizan al sistema en ejecución. Todas las métricas admiten el filtrado por intervalo de fechas, estado de la tarea y agente.

  • Totales de tareas por estado. Recuentos de tareas en cada estado (encolada, en proceso, completada, fallida, cancelada), actualizados en tiempo real a medida que ocurren las transiciones de estado.
  • Tokens de entrada y salida. Recuentos de tokens agregados y por tarea extraídos de MLflow Traces.
  • Coste de LLM. Tanto la estimación emitida por el modelo a partir de MLflow Traces (disponible a los pocos segundos de cada llamada al modelo).
  • Coste de computación. Coste de computación de Serverless Jobs atribuible a las ejecuciones de tareas del orquestador, extraído de system.billing.usage.
  • Tiempo de respuesta mediano. Calculado a partir de las tareas completadas. Se utiliza la mediana en lugar de la media para evitar la distorsión por valores atípicos de retroceso de reintento (retry-backoff) y la latencia de cola bajo saturación.
  • Confianza. Puntuaciones de confianza por documento devueltas por la capa de agentes de IA, presentadas junto con los resultados de las tareas.

Implementación

Los cambios de estado en la tabla de tareas activan eventos LISTEN/NOTIFY de Postgres. El backend mantiene una única conexión LISTEN y distribuye los eventos a través de Server-Sent Events (SSE) a los clientes del panel de control conectados. Los navegadores abren una conexión EventSource y reciben actualizaciones en vivo aproximadamente un segundo después de cada cambio de estado significativo. La implementación no requiere Redis, ni servidor WebSocket, ni bus de mensajes.

El sondeo se mantiene como un mecanismo de respaldo permanente con un intervalo predeterminado de diez segundos. Las conexiones de streaming a través de proxies de entrada en la nube pueden perder bytes de forma silenciosa sin activar eventos de error en el lado del cliente; el sondeo permanente garantiza que el dashboard se mantenga actualizado incluso en esos casos. Un indicador de la UI distingue entre los canales live (SSE activo) y polling (SSE no disponible).

Los datos del dashboard abarcan tres fuentes con diferentes características de latencia: Postgres (instantánea), la API de trace de MLflow (menos de un segundo) y consultas de SQL warehouse contra las tablas de facturación del sistema (ocasionalmente decenas de segundos). Las consultas rápidas alimentan cada ciclo de actualización; las consultas lentas se ejecutan solo tras la acción del usuario y devuelven un estado de carga de manera optimista hasta que los resultados estén disponibles.

Atribución de costos por aplicación

Las tablas de facturación del sistema de Databricks tienen alcance a nivel de cuenta: cada trabajo (job), cada llamada al modelo y cualquier otra aplicación contribuyen a las mismas filas de system.billing.usage. Sin delimitar el alcance, una tarjeta de "costo de OCR" a nivel de aplicación agregaría el uso de cada llamada al modelo en el espacio de trabajo (workspace).

La solución consiste en registrar qué ejecuciones de Databricks Jobs envió el orquestador (con un seguimiento en tasks.locked_by y task_attempts.run_id) y filtrar la consulta de facturación para ese conjunto. Un único SQL warehouse puede dar soporte a múltiples aplicaciones, y cada dashboard muestra únicamente su propio gasto.

La misma arquitectura de consulta se integra de forma natural con filtros definidos por el operador. Las cifras de costos, junto con cualquier otra métrica del dashboard, se pueden desglosar aún más por rango de fechas, estado de la tarea o agente, lo que permite responder a preguntas como "¿cuánto costaron las tareas fallidas en los últimos siete días?" o "¿cuál fue la mediana de gasto por tarea para el agente X este mes?" sin salir del dashboard.

Esto facilita el monitoreo, la asignación y el reporte de las cifras de costos.

Lakebase como el núcleo de la orquestación

Los patrones de Postgres como cola (Postgres-as-queue) están muy consolidados en la comunidad de ingeniería de datos. Lakebase proporciona las características operativas adicionales que hacen que este patrón sea viable como una arquitectura de producción en Databricks:

  • Cómputo con escalado automático. Lakebase escala las unidades de cómputo de Postgres hacia arriba y hacia abajo según la carga de trabajo, lo que permite al orquestador apoyarse en la base de datos sin tener que pagar por la capacidad máxima las 24 horas del día.
  • Autenticación rotada por OAuth. Lakebase utiliza tokens de OAuth de corta duración para la autenticación de conexiones. Los pools de conexiones actualizan los tokens automáticamente, lo que elimina las credenciales estáticas en la configuración de la aplicación y descarta los manuales de procedimientos (runbooks) de rotación.
  • Integración con Unity Catalog. Lakebase comparte identidad, permisos y gobernanza con el resto de Databricks. El principal de servicio del orquestador recibe permisos explícitos en las tablas tasks y results; no se requiere ninguna configuración de IAM independiente.
  • Ramificación y snapshots. Clonar una tabla de tareas de producción en un entorno de desarrollo para la depuración de errores (debugging) es una operación estándar de Lakebase, admitida de forma nativa.

Estas capacidades eliminan la carga operativa que normalmente motiva a los equipos a adoptar agentes de mensajes (message brokers) administrados en lugar de Postgres autohospedado para la cola de tareas.

Impacto y conclusión

En CLA, el patrón de orquestación descrito aquí respalda un flujo de trabajo de procesamiento de documentos en producción que reduce el tiempo de extracción de horas a minutos. La arquitectura utiliza servicios nativos de Databricks con Lakebase Postgres en el centro para gestionar las colas, la programación y la observabilidad sin necesidad de sistemas externos. Esto reduce la sobrecarga de integración al tiempo que aprovecha al máximo una plataforma unificada diseñada para escalar.

En producción, este patrón proporciona una gestión de tareas duradera, control de prioridades, programación adaptada a los límites de velocidad (rate-limit), visibilidad en tiempo real y seguimiento de costos por tarea. Juntas, estas capacidades ofrecen una forma práctica de orquestar cargas de trabajo de agentes (agentic workloads) manteniendo la infraestructura circundante simple.

¿Listo para simplificar la orquestación de agentes de AI en una sola plataforma? Prueba Databricks Free Edition, crea tu primer proyecto de Lakebase Postgres y tu primera Databricks App, y luego sigue la demostración de 10 minutos de MLflow Tracing para añadir observabilidad de extremo a extremo al flujo de trabajo de tus agentes.

(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.