Passa al contenuto principale
Prodotto

Semplifica l'orchestrazione degli agenti AI con Lakebase Postgres

Come CLA ha creato una soluzione nativa di Databricks per task a lunga esecuzione, osservabilità e attribuzione dei costi

di Li Yu, Michelle JanneyCoyle, Jon Cormack, Yarri Bryn, Alec Sorensen e Darshana Nair

  • Coda di task scalabile su Postgres: un'analisi approfondita dei pattern che trasformano una coppia di tabelle Lakebase in una coda durevole, concorrente e resiliente ai crash per task di agenti a lunga esecuzione, senza richiedere broker, cache o scheduler.
  • Architettura completamente nativa di Databricks: un design di riferimento che unisce Lakebase, Databricks Apps, Lakeflow Jobs, MLflow e Unity Catalog Volumes in una pipeline end-to-end per il parsing di documenti basato su agenti, senza alcuna infrastruttura esterna da gestire.
  • Osservabilità in tempo reale e invarianti: un'analisi approfondita sull'uso dei trigger LISTEN/NOTIFY di Postgres abbinati a Server-Sent Events (SSE) per creare una dashboard operatore a bassa latenza che traccia automaticamente costi e task con zero overhead.

Introduzione

Tradizionalmente, l'auditing è un processo noioso che spesso richiede una revisione dettagliata dei documenti e l'estrazione di informazioni. Per accelerare questo processo, CLA (CliftonLarsonAllen LLP), una delle principali società di servizi professionali con una crescente presenza globale, ha collaborato con il team Forward Deployed Engineering di Databricks per creare e mettere in produzione una soluzione di auditing agentica. Insieme, abbiamo sviluppato un'applicazione di elaborazione dei documenti che riduce i tempi di estrazione da ore a minuti, senza compromettere la qualità. L'applicazione è interamente basata su Databricks, utilizzando Lakebase Postgres, Databricks Apps, Lakeflow Jobs, MLflow e Unity Catalog Volumes. In questo blog ci concentreremo su un componente chiave di quel sistema: il livello di orchestrazione basato su Lakebase.

Il livello di orchestrazione è responsabile del coordinamento dei task a esecuzione prolungata, della gestione dei tentativi, dell'attribuzione dei costi e della fornitura di visibilità in tempo reale. Grazie a Lakebase e Databricks Apps, eliminiamo la necessità di un'infrastruttura separata per l'accodamento, l'orchestrazione e l'osservabilità.

Lakebase rende inoltre questa architettura pratica su scala separando lo storage dal compute. A differenza delle tradizionali distribuzioni Postgres, il compute può scalare con la domanda mentre lo storage rimane durevole e indipendente. Insieme, queste funzionalità rendono Lakebase una base pratica per un modello di orchestrazione più semplice e scalabile per carichi di lavoro agentici a esecuzione prolungata su Databricks.

Sfide di orchestrazione per i carichi di lavoro agentici

Il parsing dei documenti è un carico di lavoro agentico molto comune e ad alto volume. Le aziende di tutti i settori devono convertire grandi volumi di contratti, fatture, documenti finanziari e altri documenti in dati strutturati. L'esecuzione di questo processo su scala fa emergere cinque distinti problemi dei sistemi distribuiti:

  • Latenza imprevedibile per singolo task: una fattura di due pagine può essere elaborata in pochi secondi, mentre un contratto di duecento pagine può richiedere diversi minuti, rendendo difficile prevedere la durata dell'esecuzione di ogni singolo task.
  • Throttling sensibile ai limiti di frequenza (rate limit): gli endpoint dei modelli LLM e di visione artificiale limitano il numero di richieste e token che possono elaborare in un determinato periodo. L'invio simultaneo di centinaia di task può superare tali limiti, attivare il throttling e causare tentativi ripetuti. L'orchestratore deve limitare proattivamente il lavoro in corso (in base al numero di task simultanei, al budget di token o a entrambi) anziché affidarsi esclusivamente a retry reattivi.
  • Assegnazione delle priorità ai carichi di lavoro: gli invii urgenti non dovrebbero essere ritardati da grandi batch in blocco. La priorità per singolo task garantisce che il lavoro a priorità più elevata (invii interattivi, richieste di livello premium, rielaborazioni avviate dall'operatore) venga inviato per primo.
  • Attribuzione dei costi per task: i team finanziari devono attribuire la spesa a task, clienti e agenti specifici, suddivisi per utilizzo di token AI e consumo di risorse di calcolo (compute).
  • Visibilità del progresso in tempo reale: gli utenti che caricano centinaia di documenti hanno bisogno di una visualizzazione dello stato di avanzamento in tempo reale.

Molte organizzazioni combinano diversi sistemi specializzati per l'orchestrazione e l'osservabilità. Ciascun sistema comporta requisiti specifici in termini di infrastruttura, autenticazione, monitoraggio e operatività, oltre al lavoro necessario per integrarli. Per i task agentici indipendenti a esecuzione prolungata, tale sovraccarico è sproporzionato rispetto all'effettiva complessità di pianificazione.

La soluzione nativa di Databricks che abbiamo sviluppato soddisfa tutti i requisiti di cui sopra, con Lakebase come base.

Architettura della soluzione

Architettura della soluzione

L'intero stack applicativo è composto esclusivamente da servizi Databricks:

  • Applicazione Web (Databricks Apps). Un'interfaccia utente basata su FastAPI in cui gli utenti caricano file PDF (archiviati in Unity Catalog Volumes) e inviano richieste di parsing. Le richieste vengono scritte direttamente nella tabella dei task di Lakebase.
  • Lakebase. Un database Postgres con scalabilità automatica che ospita lo stato relazionale dell'orchestratore in tabelle correlate: tasks (documenti da analizzare, contenente stato, informazioni sul lease e il risultato strutturato) e task_attempts (una riga per tentativo di esecuzione, che acquisisce l'ID di esecuzione del Job Databricks, l'ID di tracciamento di MLflow e i metadati dei costi per tentativo). Lakebase funge da unica fonte di verità (single source of truth) per lo stato dell'orchestratore.
  • Orchestratore (Databricks Apps). Un daemon worker a esecuzione prolungata e una dashboard operatore. Il daemon rimuove i task dalla coda (dequeue) di Lakebase, li invia al livello degli agenti AI e riscrive i risultati. La dashboard legge le stesse tabelle per mostrare lo stato in tempo reale.
  • Agenti AI (Lakeflow Jobs). I Lakeflow Jobs eseguono il lavoro di parsing. Ciascun Job legge un PDF da Unity Catalog Volumes, lo elabora tramite Intelligent Document Processing e chiamate a modelli LLM/visione artificiale, memorizza l'output analizzato in Lakebase e invia una notifica all'orchestratore tramite webhook. MLflow Tracing acquisisce i dettagli di esecuzione come le chiamate al modello, l'utilizzo dei token, la latenza e i metadati dei costi.

I dati fluiscono tra i componenti come segue. L'applicazione Web scrive i PDF in Unity Catalog Volumes e le richieste di parsing in Lakebase. L'orchestratore rimuove i task dalla coda di Lakebase e invia i Job Databricks al livello degli agenti AI. Gli agenti AI elaborano i documenti, riscrivono i risultati in Lakebase e inviano un feedback all'orchestratore con gli aggiornamenti di stato.

Grazie a queste funzionalità integrate su Databricks, non abbiamo dovuto fare affidamento su message broker esterni (Kafka, Redis), scheduler separati (Airflow, Temporal) o livelli di caching dedicati.

Implementazione della coda dei task

La coda dei task è supportata da due tabelle Postgres in Lakebase. La tabella tasks contiene una riga per unità logica di lavoro, registrando lo stato corrente del task, le informazioni sul lease, l'assegnazione dell'agente, l'estrazione principale (parent) e il risultato finale. La tabella task_attempts contiene una riga per tentativo di esecuzione, acquisendo l'ID di esecuzione del Job Databricks, l'ID di tracciamento di MLflow e i metadati dei costi per tentativo. La relazione padre-figlio supporta i tentativi (un singolo task può avere più tentativi) e preserva l'osservabilità a livello di tentativo per l'attribuzione dei costi e il debug.

Una coppia di tabelle Postgres da sola non costituisce ancora una coda di task. Quattro pattern nativi di Postgres le trasformano in una coda robusta, simultanea, resiliente ai crash e sensibile ai limiti di frequenza, adatta per carichi di lavoro agentici a esecuzione prolungata.

Rimozione simultanea dalla coda sensibile alle priorità

Una query di dequeue di base potrebbe selezionare il task successivo disponibile utilizzando WHERE status = 'enqueued' and LIMIT batch_size. Sebbene tale query identifichi correttamente un task in coda, non è sufficiente quando più worker eseguono il dequeue simultaneamente. Senza il blocco delle righe (row locking), diversi worker potrebbero selezionare lo stesso task prima che lo stato venga aggiornato.

L'aggiunta di FOR UPDATE SKIP LOCKED rende il dequeue sicuro per la concorrenza. Ogni worker blocca la riga selezionata, mentre gli altri worker saltano quella riga e passano al task successivo disponibile. Inoltre, una clausola ORDER BY priority DESC, created_at garantisce che i task con priorità più alta vengano selezionati per primi, preservando l'ordinamento FIFO all'interno di ciascun livello di priorità.

L'istruzione completa, sicura per la concorrenza e stabile per priorità, è:

Ripristino dai crash tramite blocco basato su lease

I worker possono interrompersi a metà di un task a causa dell'eviction delle VM, di condizioni di memoria esaurita (out-of-memory) o di eventi di deployment. Se i task interrotti rimangono contrassegnati come in fase di elaborazione, verrebbero trattenuti a tempo indefinito. La soluzione consiste nel registrare un lease in scadenza al momento del dequeue:

Uno sweeper periodico reinserisce in coda qualsiasi task il cui lease_expires_at sia scaduto. I task trattenuti dai worker interrotti vengono recuperati automaticamente in pochi minuti, senza richiedere un servizio di coordinamento esterno.

Throttling sensibile ai rate limit

Gli endpoint dei modelli LLM e di visione solitamente applicano due quote distinte: un limite di richieste al secondo e un limite di token al minuto (TPM). Una singola strategia di throttling raramente si adatta a entrambi. L'orchestratore supporta tre modalità, selezionate per ciascun agente tramite configurazione.

Limite di concorrenza. Un parametro MAX_CONCURRENT_TASKS limita il numero di task che l'orchestratore distribuisce contemporaneamente. Il limite viene applicato al momento del dequeue contando le righe PROCESSING correnti nella tabella dei task:

Se il conteggio è pari o superiore al limite, non viene rimosso dalla coda alcun nuovo task. Ancorare il controllo al conteggio delle righe del database, anziché alla dimensione della coda dell'esecutore locale, mantiene accurato il limite tra i riavvii dei worker, i recuperi dei lease e i deployment multi-replica. Questa modalità è ideale per gli endpoint vincolati da limiti di richieste al secondo in cui l'utilizzo dei token di ciascun task è approssimativamente uniforme.

Budget di token. Un parametro MAX_TPM limita la frequenza di token stimata per i task in corso. L'orchestratore stima il conteggio dei token di un task e somma la frequenza di token stimata per tutti i task PROCESSING. Un nuovo task viene rimosso dalla coda solo se la somma più i token stimati del nuovo task rientra nel budget.

Limite combinato. Quando sono configurati sia MAX_CONCURRENT_TASKS sia MAX_TPM, l'orchestratore applica il vincolo più restrittivo. Questa modalità gestisce carichi di lavoro che sono limitati dalla concorrenza in un regime (molti task brevi ed economici) e limitati dai token in un altro (un singolo documento molto lungo che satura la quota al minuto).

In tutte e tre le modalità, la decisione di throttling viene presa al momento del dequeue all'interno della stessa transazione di FOR UPDATE SKIP LOCKED. Un task che non rientra nella quota corrente rimane in coda e viene riconsiderato nel ciclo di dequeue successivo: nessun stato di pianificazione separato, nessuna coda di attesa in memoria, nessun livello di coordinamento tra le repliche dei worker.

Callback webhook idempotenti

Quando il livello AI Agents completa un task, invia una callback all'orchestratore con il risultato. La consegna della callback non è exactly-once: Databricks potrebbe riprovare, le reti potrebbero interrompersi e i proxy potrebbero riconsegnare. L'handler della callback è progettato per essere idempotente: accetta sia lo stato PROCESSING che ENQUEUED e tratta i task già terminati come no-op. Payload identici producono risultati identici, eliminando il rischio di doppia fatturazione o di elaborazioni duplicate.

Questi quattro pattern combinati producono una coda di task corretta in condizioni di concorrenza, durevole in caso di crash, sensibile ai rate limit sotto carico e idempotente in caso di retry. La visibilità in tempo reale sul sistema in esecuzione è fornita da un meccanismo separato descritto nella sezione successiva.

Dashboard operatore in tempo reale

Quando molti documenti sono in corso di elaborazione, gli operatori hanno bisogno di una visione chiara delle prestazioni degli agenti, dello stato dei task e del costo del carico di lavoro. Non dovrebbero dover interrogare continuamente l'orchestratore o affidarsi a una piattaforma di metriche separata. L'orchestratore integra questa funzionalità direttamente in una singola dashboard esposta dalla stessa Databricks App che esegue il daemon del worker.

Funzionalità della dashboard

La dashboard rivolta agli operatori mostra un insieme di metriche operative che caratterizzano il sistema in esecuzione. Tutte le metriche supportano il filtraggio per intervallo di date, stato del task e agente.

  • Totali dei task per stato. Conteggi dei task in ciascuno stato (enqueued, processing, completed, failed, cancelled), aggiornati in tempo reale al verificarsi delle transizioni di stato.
  • Token di input e output. Conteggi di token per task e aggregati estratti da MLflow Traces.
  • Costo LLM. Sia la stima emessa dal modello da MLflow Traces (disponibile entro pochi secondi da ciascuna chiamata al modello).
  • Costo di calcolo. Costo di calcolo dei Serverless Jobs attribuibile alle esecuzioni dei task dell'orchestratore, estratto da system.billing.usage.
  • Tempo di risposta mediano. Calcolato sui task completati. La mediana viene utilizzata al posto della media per evitare distorsioni dovute a outlier di retry-backoff e alla latenza di coda in caso di saturazione.
  • Confidenza. Punteggi di confidenza per documento restituiti dal livello AI Agents, mostrati insieme ai risultati dei task.

Implementazione

Le modifiche di stato nella tabella dei task attivano eventi Postgres LISTEN/NOTIFY. Il backend mantiene una singola connessione LISTEN e distribuisce gli eventi tramite Server-Sent Events (SSE) ai client della dashboard connessi. I browser aprono una connessione EventSource e ricevono aggiornamenti in tempo reale entro circa un secondo da ogni modifica di stato significativa. L'implementazione non richiede Redis, nessun server WebSocket e nessun bus di messaggi.

Il polling viene mantenuto come fallback permanente a un intervallo predefinito di dieci secondi. Le connessioni in streaming tramite proxy di ingresso cloud possono perdere byte in modo silenzioso senza attivare eventi di errore lato client; il polling permanente garantisce che la dashboard rimanga aggiornata anche in questi casi. Un indicatore UI distingue tra i canali live (SSE attivo) e polling (SSE non disponibile).

I dati della dashboard provengono da tre fonti con diverse caratteristiche di latenza: Postgres (istantanea), l'API di tracciamento di MLflow (inferiore al secondo) e le query del warehouse sulle tabelle di fatturazione del sistema (occasionalmente nell'ordine di decine di secondi). Le query veloci alimentano ogni ciclo di aggiornamento; le query lente vengono eseguite solo su azione dell'utente e restituiscono ottimisticamente uno stato di caricamento fino a quando i risultati non sono disponibili.

Attribuzione dei costi per applicazione

Le tabelle di fatturazione del sistema Databricks hanno come ambito l'intero account: ogni job, ogni chiamata al modello e ogni altra applicazione contribuiscono alle stesse righe di system.billing.usage. Senza questa limitazione dell'ambito, un riquadro "costo OCR" a livello di applicazione aggregherebbe l'utilizzo di ogni chiamata al modello nell'area di lavoro.

La soluzione consiste nel registrare quali esecuzioni di Databricks Job sono state inviate dall'orchestratore (tracciate in tasks.locked_by e task_attempts.run_id) e filtrare la query di fatturazione su quel set. Un singolo SQL warehouse può supportare più applicazioni e ogni dashboard mostra solo la propria spesa.

La stessa architettura di query si integra naturalmente con i filtri guidati dall'operatore. Le cifre dei costi, insieme a ogni altra metrica della dashboard, possono essere ulteriormente limitate per intervallo di date, stato del task o agente, supportando domande come "quanto sono costati i task non riusciti negli ultimi sette giorni?" o "qual è stata la spesa mediana per task per l'agente X questo mese?" senza uscire dalla dashboard.

Questo rende le cifre dei costi facili da monitorare, allocare e segnalare.

Lakebase come spina dorsale dell'orchestrazione

I pattern Postgres-as-queue sono ormai consolidati nella community del data engineering. Lakebase fornisce le caratteristiche operative aggiuntive che rendono questo pattern praticabile come architettura di produzione su Databricks:

  • Calcolo con scalabilità automatica. Lakebase scala le unità di calcolo di Postgres verso l'alto o verso il basso in base al carico di lavoro, consentendo all'orchestratore di appoggiarsi al database senza pagare per la capacità di picco 24 ore su 24.
  • Autenticazione con rotazione OAuth. Lakebase utilizza token OAuth a breve durata per l'autenticazione delle connessioni. I pool di connessioni aggiornano i token automaticamente, eliminando le credenziali statiche nella configurazione dell'applicazione ed evitando i runbook di rotazione.
  • Integrazione con Unity Catalog. Lakebase condivide identità, autorizzazioni e governance con il resto di Databricks. Il service principal dell'orchestratore riceve autorizzazioni esplicite sulle tabelle tasks e results tabelle; non è richiesta alcuna configurazione IAM separata.
  • Branching e snapshot. La clonazione di una tabella di task di produzione in un ambiente di sviluppo per il debug è un'operazione standard di Lakebase, supportata nativamente.

Queste funzionalità eliminano il carico operativo che in genere spinge i team ad adottare message broker gestiti al posto di Postgres self-hosted per l'accodamento dei task.

Impatto e conclusioni

Presso CLA, il pattern di orchestrazione qui descritto supporta un flusso di lavoro di elaborazione dei documenti in produzione che riduce i tempi di estrazione da ore a minuti. L'architettura utilizza servizi nativi di Databricks con Lakebase Postgres al centro per gestire l'accodamento, la pianificazione e l'osservabilità senza richiedere sistemi esterni. Questo riduce il sovraccarico di integrazione sfruttando appieno una piattaforma unificata progettata per scalare.

In produzione, questo pattern offre una gestione affidabile dei task, il controllo delle priorità, una pianificazione sensibile ai limiti di velocità (rate-limit), visibilità in tempo reale e il tracciamento dei costi per singolo task. Insieme, queste funzionalità offrono un modo pratico per orchestrare flussi di lavoro agentici mantenendo semplice l'infrastruttura circostante.

Vuoi semplificare l'orchestrazione degli agenti AI su un'unica piattaforma? Prova Databricks Free Edition, crea il tuo primo progetto Lakebase Postgres e la tua prima Databricks App, quindi segui la demo di 10 minuti di MLflow Tracing per aggiungere l'osservabilità end-to-end al flusso di lavoro del tuo agente.

(Questo post sul blog è stato tradotto utilizzando strumenti basati sull'intelligenza artificiale) Post originale

Ricevi gli ultimi articoli nella tua casella di posta

Iscriviti al nostro blog e ricevi gli ultimi articoli direttamente nella tua casella di posta.