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

El stack completo de la aplicación consta exclusivamente de servicios de Databricks:
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.
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.
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:
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.
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.
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.
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.
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.
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.
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.
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:
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.
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
Suscríbete a nuestro blog y recibe las últimas publicaciones directamente en tu bandeja de entrada.