Revenir au contenu principal
Produit

Simplifiez l'orchestration d'agents AI avec Lakebase Postgres

Comment CLA a conçu une solution native de Databricks pour les tâches de longue durée, l'observabilité et l'imputation des coûts

par Li Yu, Michelle JanneyCoyle, Jon Cormack, Yarri Bryn, Alec Sorensen et Darshana Nair

  • File d'attente de tâches prête pour la mise à l'échelle sur Postgres : une plongée au cœur des modèles qui transforment une paire de tables Lakebase en une file d'attente durable, concurrente et résiliente aux pannes pour les tâches d'agent de longue durée — sans qu'aucun courtier, cache ou planificateur ne soit requis.
  • Architecture entièrement native de Databricks : une conception de référence qui assemble Lakebase, Databricks Apps, Lakeflow Jobs, MLflow et Unity Catalog Volumes en un pipeline de bout en bout pour l'analyse documentaire basée sur des agents, sans infrastructure externe à gérer.
  • Observabilité en temps réel et invariants : une plongée au cœur de l'utilisation des déclencheurs Postgres LISTEN/NOTIFY associés aux Server-Sent Events (SSE) pour créer un tableau de bord opérateur à faible latence qui suit automatiquement les coûts et les tâches avec un surdébit nul.

Introduction

Traditionnellement, l'audit est un processus fastidieux qui nécessite souvent un examen détaillé des documents et une extraction d'informations. Pour accélérer ce processus, CLA (CliftonLarsonAllen LLP), un cabinet de services professionnels de premier plan ayant une présence mondiale croissante, a collaboré avec l'équipe Forward Deployed Engineering de Databricks pour concevoir et mettre en production une solution d'audit agentique. Ensemble, nous avons développé une application de traitement de documents qui réduit le temps d'extraction d'heures en minutes, sans aucun compromis sur la qualité. L'application est entièrement construite sur Databricks, en utilisant Lakebase Postgres, Databricks Apps, Lakeflow Jobs, MLflow et Unity Catalog Volumes. Dans ce blog, nous nous concentrons sur un composant clé de ce système : la couche d'orchestration optimisée par Lakebase.

La couche d'orchestration est chargée de coordonner les tâches de longue durée, de gérer les tentatives, d'attribuer les coûts et de fournir une visibilité en temps réel. Avec Lakebase et Databricks Apps, nous éliminons le besoin d'une infrastructure distincte pour la mise en file d'attente, l'orchestration et l'observabilité.

Lakebase rend également cette architecture pratique à grande échelle en séparant le stockage du calcul. Contrairement aux déploiements Postgres traditionnels, le calcul peut s'adapter à la demande tandis que le stockage reste durable et indépendant. Ensemble, ces capacités font de Lakebase une base pratique pour un modèle d'orchestration plus simple et évolutif pour les charges de travail agentiques de longue durée sur Databricks.

Défis d'orchestration pour les charges de travail agentiques

L'analyse de documents est une charge de travail agentique très courante et à volume élevé. Les entreprises de tous les secteurs doivent convertir de grands volumes de contrats, de factures, de documents financiers et d'autres documents en données structurées. L'exécution de cette tâche à grande échelle fait apparaître cinq problèmes distincts liés aux systèmes distribués :

  • Latence imprévisible par tâche : une facture de deux pages peut être traitée en quelques secondes, tandis qu'un contrat de deux cents pages peut prendre plusieurs minutes, ce qui rend difficile la prévision de la durée d'exécution d'une tâche individuelle.
  • Régulation prenant en compte les limites de débit : les points de terminaison des modèles LLM et de vision limitent le nombre de requêtes et de jetons qu'ils peuvent traiter sur une période donnée. L'envoi de centaines de tâches à la fois peut dépasser ces limites, déclencher une régulation et entraîner des tentatives répétées. L'orchestrateur doit limiter de manière proactive le travail en cours (par le nombre de tâches simultanées, par le budget de jetons, ou les deux) plutôt que de s'en remettre uniquement à des tentatives réactives.
  • Priorisation de la charge de travail : les soumissions urgentes ne doivent pas être retardées par des lots volumineux. La priorité par tâche garantit que le travail de priorité supérieure (soumissions interactives, demandes de niveau premium, retraitements initiés par l'opérateur) est distribué en premier.
  • Attribution des coûts par tâche : les équipes financières doivent attribuer les dépenses à des tâches, des clients et des agents spécifiques, réparties par utilisation de jetons d'IA et consommation de calcul.
  • Visibilité de la progression en temps réel : les utilisateurs qui importent des centaines de documents ont besoin d'une vue de la progression en direct.

De nombreuses organisations combinent plusieurs systèmes spécialisés pour l'orchestration et l'observabilité. Chaque système apporte sa propre infrastructure, ses propres exigences d'authentification, de surveillance et de fonctionnement, ainsi que le travail requis pour les intégrer. Pour les tâches agentiques indépendantes et de longue durée, cette surcharge est disproportionnée par rapport à la complexité réelle de la planification.

La solution native Databricks que nous avons développée répond à toutes les exigences ci-dessus avec Lakebase comme fondation.

Architecture de la solution

Architecture de la solution

L'ensemble de la pile d'applications se compose exclusivement de services Databricks :

  • Application Web (Databricks Apps). Une interface utilisateur basée sur FastAPI où les utilisateurs importent des PDF (stockés dans Unity Catalog Volumes) et soumettent des demandes d'analyse. Les requêtes sont écrites directement dans la table des tâches Lakebase.
  • Lakebase. Une base de données Postgres à mise à l'échelle automatique (autoscaling) hébergeant l'état relationnel de l'orchestrateur dans des tables associées : tasks (documents à analyser, contenant le statut, les informations de bail et le résultat structuré), et task_attempts (une ligne par tentative d'exécution, capturant l'ID d'exécution du Job Databricks, l'ID de trace MLflow et les métadonnées de coût par tentative). Lakebase sert de source unique de vérité pour l'état de l'orchestrateur.
  • Orchestrateur (Databricks Apps). Un démon worker de longue durée et un tableau de bord opérateur. Le démon retire les tâches de la file d'attente de Lakebase, les distribue à la couche des agents d'IA et réécrit les résultats. Le tableau de bord lit les mêmes tables pour afficher le statut en temps réel.
  • Agents d'IA (Lakeflow Jobs). Les Lakeflow Jobs exécutent le travail d'analyse. Chaque Job lit un PDF à partir de Unity Catalog Volumes, le traite via le traitement intelligent de documents et des appels de modèle de vision/LLM, stocke la sortie analysée dans Lakebase et en accuse réception à l'orchestrateur via un webhook. MLflow Tracing capture les détails d'exécution tels que les appels de modèle, l'utilisation des jetons, la latence et les métadonnées de coût.

Les données circulent entre les composants de la manière suivante. L'application Web écrit les PDF dans Unity Catalog Volumes et les requêtes d'analyse dans Lakebase. L'orchestrateur retire les tâches de la file d'attente de Lakebase et distribue les Jobs Databricks à la couche des agents d'IA. Les agents d'IA traitent les documents, réécrivent les résultats dans Lakebase et rappellent l'orchestrateur avec des mises à jour de statut.

Grâce à ces fonctionnalités intégrées sur Databricks, nous n'avons pas eu besoin de nous appuyer sur des courtiers de messages externes (Kafka, Redis), des planificateurs distincts (Airflow, Temporal) ou des couches de mise en cache dédiées.

Implémentation de la file d'attente des tâches

La file d'attente des tâches s'appuie sur deux tables Postgres dans Lakebase. La table tasks contient une ligne par unité logique de travail, enregistrant le statut actuel de la tâche, les informations de bail, l'attribution de l'agent, l'extraction parente et le résultat final. La table task_attempts contient une ligne par tentative d'exécution, capturant l'ID d'exécution du Job Databricks, l'ID de trace MLflow et les métadonnées de coût par tentative. La relation parent-enfant prend en charge les tentatives (une seule tâche peut faire l'objet de plusieurs tentatives) et préserve l'observabilité au niveau de la tentative pour l'attribution des coûts et le débogage.

Deux tables Postgres ne constituent pas à elles seules une file d'attente de tâches. Quatre modèles natifs Postgres les transforment en une file d'attente robuste, simultanée, résiliente aux pannes et prenant en compte les limites de débit, adaptée aux charges de travail agentiques de longue durée.

Retrait de la file d'attente simultané et prenant en compte la priorité

Une requête de retrait de file d'attente de base peut sélectionner la tâche disponible suivante en utilisant WHERE status = 'enqueued' and LIMIT batch_size. Bien que cette requête identifie correctement une tâche mise en file d'attente, elle n'est pas suffisante lorsque plusieurs workers effectuent des retraits de file d'attente de manière simultanée. Sans verrouillage de ligne, plusieurs workers peuvent sélectionner la même tâche avant que le statut ne soit mis à jour.

L'ajout de FOR UPDATE SKIP LOCKED rend le retrait de la file (dequeue) sûr pour la concurrence. Chaque worker verrouille la ligne qu'il sélectionne, tandis que les autres workers ignorent cette ligne et passent à la tâche disponible suivante. De plus, une clause ORDER BY priority DESC, created_at garantit que les tâches prioritaires sont sélectionnées en premier, tout en préservant l'ordre FIFO au sein de chaque niveau de priorité.

L'instruction complète, sûre pour la concurrence et stable en termes de priorité, est :

Récupération après plantage via un verrouillage basé sur des baux

Les workers peuvent s'arrêter en cours de tâche en raison d'une éviction de VM, de conditions de manque de mémoire (out-of-memory) ou d'événements de déploiement. Si les tâches interrompues restent marquées comme étant en cours de traitement, elles seraient bloquées indéfiniment. La solution consiste à enregistrer un bail expirant au moment du retrait de la file :

Un nettoyeur (sweeper) périodique réenfile toute tâche dont le lease_expires_at est dépassé. Les tâches détenues par des workers arrêtés sont récupérées automatiquement en quelques minutes, sans nécessiter de service de coordination externe.

Régulation du débit prenant en compte les limites de taux

Les points de terminaison (endpoints) des modèles LLM et de vision imposent généralement deux quotas distincts : une limite de requêtes par seconde et une limite de jetons par minute (TPM). Une seule stratégie de régulation convient rarement aux deux. L'orchestrateur prend en charge trois modes, sélectionnés par agent via la configuration.

Limite de concurrence. Un paramètre MAX_CONCURRENT_TASKS limite le nombre de tâches que l'orchestrateur distribue simultanément. La limite est appliquée au moment du retrait de la file en comptant les lignes PROCESSING actuelles dans la table des tâches :

Si le nombre est égal ou supérieur à la limite, aucune nouvelle tâche n'est retirée de la file. L'ancrage de cette vérification sur le nombre de lignes de la base de données, plutôt que sur la taille de la file d'attente de l'exécuteur local, permet de maintenir une limite précise malgré les redémarrages de workers, les récupérations de baux et les déploiements multi-répliques. Ce mode est particulièrement adapté aux points de terminaison limités par le nombre de requêtes par seconde, où l'utilisation des jetons pour chaque tâche est relativement uniforme.

Budget de jetons. Un paramètre MAX_TPM limite le taux de jetons projeté pour les tâches en cours. L'orchestrateur estime le nombre de jetons d'une tâche et additionne le taux de jetons projeté pour toutes les tâches PROCESSING. Une nouvelle tâche n'est retirée de la file que si cette somme, augmentée des jetons projetés de la nouvelle tâche, respecte le budget.

Limite combinée. Lorsque MAX_CONCURRENT_TASKS et MAX_TPM sont tous deux configurés, l'orchestrateur applique la contrainte la plus stricte. Ce mode gère les charges de travail qui sont limitées par la concurrence dans un cas (nombreuses tâches courtes et peu coûteuses) et limitées par les jetons dans un autre (un seul document très long saturant le quota par minute).

Dans les trois modes, la décision de régulation est prise au moment du retrait de la file, au sein de la même transaction que FOR UPDATE SKIP LOCKED. Une tâche qui ne peut pas s'intégrer dans le quota actuel reste dans la file d'attente et est réexaminée lors du cycle de retrait suivant — pas d'état de planification distinct, pas de file d'attente en mémoire, pas de couche de coordination entre les répliques de workers.

Callbacks Webhook Idempotents

Lorsque la couche AI Agents termine une tâche, elle envoie un callback à l'orchestrateur avec le résultat. La distribution des callbacks n'est pas strictement unique (exactly-once) : Databricks peut effectuer des tentatives, les réseaux peuvent subir des interruptions et les proxys peuvent renvoyer les messages. Le gestionnaire de callbacks est conçu pour être idempotent : il accepte à la fois les états PROCESSING et ENQUEUED, et traite les tâches déjà terminées comme des opérations sans effet (no-ops). Des charges utiles (payloads) identiques produisent des résultats identiques, éliminant ainsi le risque de double facturation ou de double traitement.

La combinaison de ces quatre modèles permet d'obtenir une file d'attente de tâches correcte en situation de concurrence, durable face aux plantages, consciente des limites de taux sous la charge, et idempotente lors des tentatives. Une visibilité en temps réel sur le système en cours d'exécution est assurée par un mécanisme distinct décrit dans la section suivante.

Tableau de bord opérateur en temps réel

Lorsque de nombreux documents sont en cours de traitement, les opérateurs ont besoin d'une vue claire des performances des agents, du statut des tâches et du coût de la charge de travail. Ils ne devraient pas avoir à interroger continuellement l'orchestrateur ni à dépendre d'une plateforme de métriques distincte. L'orchestrateur intègre cette fonctionnalité directement dans un tableau de bord unique, affiché par la même Databricks App qui exécute le démon du worker.

Fonctionnalités du tableau de bord

Le tableau de bord destiné aux opérateurs présente un ensemble de métriques opérationnelles qui caractérisent le système en cours d'exécution. Toutes les métriques prennent en charge le filtrage par plage de dates, statut de tâche et agent.

  • Nombre total de tâches par statut. Nombre de tâches dans chaque état (en attente, en cours de traitement, terminées, échouées, annulées), mis à jour en temps réel au fur et à mesure des transitions d'état.
  • Jetons d'entrée et de sortie. Nombre de jetons par tâche et agrégés, issus de MLflow Traces.
  • Coût LLM. L'estimation émise par le modèle à partir de MLflow Traces (disponible quelques secondes après chaque appel de modèle).
  • Coût de calcul. Coût de calcul des Serverless Jobs attribuable aux exécutions de tâches de l'orchestrateur, extrait de system.billing.usage.
  • Temps de réponse médian. Calculé sur l'ensemble des tâches terminées. La médiane est utilisée à la place de la moyenne pour éviter toute distorsion due aux valeurs aberrantes de délai d'attente (retry-backoff) et à la latence en fin de file d'attente en cas de saturation.
  • Confiance. Scores de confiance par document renvoyés par la couche AI Agents, affichés aux côtés des résultats des tâches.

Implémentation

Les changements d'état dans la table des tâches déclenchent des événements Postgres LISTEN/NOTIFY. Le backend maintient une connexion LISTEN unique et distribue les événements via Server-Sent Events (SSE) aux clients du tableau de bord connectés. Les navigateurs ouvrent une connexion EventSource et reçoivent des mises à jour en direct dans la seconde environ suivant chaque changement d'état significatif. L'implémentation ne nécessite aucun serveur Redis, aucun serveur WebSocket, ni aucun bus de messages.

Le polling est conservé comme solution de repli permanente avec un intervalle par défaut de dix secondes. Les connexions de streaming via des proxys d'entrée cloud peuvent perdre des octets de manière silencieuse sans déclencher d'événements d'erreur côté client ; le polling permanent garantit que le tableau de bord reste à jour même dans de tels cas. Un indicateur de l'UI permet de distinguer les canaux live (SSE actif) et polling (SSE indisponible).

Les données du tableau de bord proviennent de trois sources aux caractéristiques de latence différentes : Postgres (instantané), l'API de trace de MLflow (inférieure à la seconde) et les requêtes d'entrepôt de données sur les tables de facturation système (parfois de l'ordre de quelques dizaines de secondes). Les requêtes rapides alimentent chaque cycle de rafraîchissement ; les requêtes lentes s'exécutent uniquement lors d'une action de l'utilisateur et renvoient de manière optimiste un état de chargement jusqu'à ce que les résultats soient disponibles.

Attribution des coûts par application

Les tables de facturation système de Databricks sont définies au niveau du compte : chaque tâche, chaque appel de modèle et chaque autre application contribuent aux mêmes lignes system.billing.usage. Sans cette délimitation de portée, une tuile « coût OCR » au niveau de l'application agrégerait l'utilisation de tous les appels de modèle de l'espace de travail.

La solution consiste à enregistrer les exécutions de Databricks Job soumises par l'orchestrateur (suivies dans tasks.locked_by et task_attempts.run_id) et à filtrer la requête de facturation sur cet ensemble. Un seul SQL warehouse peut prendre en charge plusieurs applications, et chaque tableau de bord ne présente que ses propres dépenses.

Cette même architecture de requête s'intègre naturellement avec des filtres définis par l'opérateur. Les chiffres de coûts, ainsi que toutes les autres métriques du tableau de bord, peuvent être affinés par plage de dates, statut de tâche ou agent, permettant de répondre à des questions telles que « quel a été le coût des tâches ayant échoué au cours des sept derniers jours ? » ou « quelle a été la dépense médiane par tâche pour l'agent X ce mois-ci ? » sans quitter le tableau de bord.

Cela permet de surveiller, d'attribuer et de signaler facilement les chiffres de coûts.

Lakebase comme pilier de l'orchestration

Les modèles de type « Postgres-as-queue » sont bien établis dans la communauté de l'ingénierie des données. Lakebase apporte les caractéristiques opérationnelles supplémentaires qui rendent ce modèle viable en tant qu'architecture de production sur Databricks :

  • Mise à l'échelle automatique du calcul. Lakebase ajuste à la hausse ou à la baisse les unités de calcul Postgres en fonction de la charge de travail, ce qui permet à l'orchestrateur de s'appuyer sur la base de données sans payer pour une capacité maximale en continu.
  • Authentification par rotation OAuth. Lakebase utilise des jetons OAuth à courte durée de vie pour l'authentification des connexions. Les pools de connexions actualisent automatiquement les jetons, éliminant ainsi les identifiants statiques dans la configuration de l'application et supprimant les guides de rotation (runbooks).
  • Intégration avec Unity Catalog. Lakebase partage l'identité, les autorisations et la gouvernance avec le reste de Databricks. Le principal de service de l'orchestrateur reçoit des autorisations explicites sur les tables tasks et results ; aucune configuration IAM distincte n'est requise.
  • Création de branches et d'instantanés. Le clonage d'une table de tâches de production dans un environnement de développement pour le débogage est une opération Lakebase standard, prise en charge nativement.

Ces fonctionnalités éliminent la lourdeur opérationnelle qui pousse généralement les équipes à adopter des courtiers de messages gérés au lieu d'un Postgres auto-hébergé pour la mise en file d'attente des tâches.

Impact et conclusion

Chez CLA, le modèle d'orchestration décrit ici prend en charge un flux de travail de traitement de documents en production qui réduit le temps d'extraction de plusieurs heures à quelques minutes. L'architecture utilise des services natifs de Databricks avec Lakebase Postgres au centre pour gérer la mise en file d'attente, la planification et l'observabilité sans nécessiter de systèmes externes. Cela réduit les coûts d'intégration tout en tirant pleinement parti d'une plateforme unifiée conçue pour évoluer.

En production, ce modèle offre une gestion durable des tâches, un contrôle des priorités, une planification prenant en compte les limites de débit, une visibilité en temps réel et un suivi des coûts par tâche. Ensemble, ces fonctionnalités offrent un moyen pratique d'orchestrer des charges de travail d'agents tout en maintenant la simplicité de l'infrastructure environnante.

Prêt à simplifier l'orchestration des agents d'AI sur une seule plateforme ? Essayez Databricks Free Edition, créez votre premier projet Lakebase Postgres et votre première Databricks App, puis suivez la démo de 10 minutes sur MLflow Tracing pour ajouter une observabilité de bout en bout à votre flux de travail d'agent.

(Cet article de blog a été traduit à l'aide d'outils basés sur l'intelligence artificielle) Article original

Recevez les derniers articles dans votre boîte mail

Abonnez-vous à notre blog et recevez les derniers articles directement dans votre boîte mail.