Como a CLA criou uma solução nativa do Databricks para tarefas de longa duração, observabilidade e atribuição de custos
por Li Yu, Michelle JanneyCoyle, Jon Cormack, Yarri Bryn, Alec Sorensen e Darshana Nair
Tradicionalmente, a auditoria é um processo tedioso que frequentemente exige uma revisão detalhada de documentos e extração de informações. Para acelerar esse processo, a CLA (CliftonLarsonAllen LLP), uma empresa líder em serviços profissionais com uma presença global crescente, trabalhou com a equipe de Forward Deployed Engineering da Databricks para criar e colocar em produção uma solução de auditoria agêntica. Juntos, desenvolvemos um aplicativo de processamento de documentos que reduz o tempo de extração de horas para minutos, sem comprometer a qualidade. O aplicativo foi desenvolvido inteiramente na Databricks, usando Lakebase Postgres, Databricks Apps, Lakeflow Jobs, MLflow e Unity Catalog Volumes. Neste blog, focamos em um componente essencial desse sistema: a camada de orquestração baseada em Lakebase.
A camada de orquestração é responsável por coordenar tarefas de longa duração, gerenciar tentativas, atribuir custos e fornecer visibilidade em tempo real. Com o Lakebase e o Databricks Apps, eliminamos a necessidade de uma infraestrutura separada para enfileiramento, orquestração e observabilidade.
O Lakebase também torna essa arquitetura prática em escala ao separar o armazenamento da computação. Ao contrário das implantações tradicionais do Postgres, a computação pode ser dimensionada de acordo com a demanda, enquanto o armazenamento permanece durável e independente. Juntos, esses recursos tornam o Lakebase uma base prática para um padrão de orquestração mais simples e escalável para cargas de trabalho agênticas de longa duração na Databricks.
A análise de documentos é uma carga de trabalho agêntica muito comum e de alto volume. Empresas de diversos setores precisam converter grandes volumes de contratos, faturas, relatórios financeiros e outros documentos em dados estruturados. Executar isso em escala traz à tona cinco problemas distintos de sistemas distribuídos:
Muitas organizações combinam vários sistemas especializados para orquestração e observabilidade. Cada sistema traz seus próprios requisitos de infraestrutura, autenticação, monitoramento e operação, além do trabalho necessário para integrá-los. Para tarefas agênticas independentes e de longa duração, essa sobrecarga é desproporcional à complexidade real do agendamento.
A solução nativa da Databricks que desenvolvemos atende a todos os requisitos acima, tendo o Lakebase como base.

A pilha completa de aplicativos consiste exclusivamente em serviços da Databricks:
O fluxo de dados entre os componentes ocorre da seguinte forma. O Web App grava PDFs no Unity Catalog Volumes e envia solicitações de análise para o Lakebase. O Orquestrador desenfileira as tarefas do Lakebase e envia Databricks Jobs para a camada de AI Agents. Os AI Agents processam os documentos, gravam os resultados de volta no Lakebase e retornam a chamada para o Orquestrador com atualizações de status.
Devido a esses recursos integrados na Databricks, não precisamos depender de agentes de mensagens externos (Kafka, Redis), agendadores separados (Airflow, Temporal) ou camadas de cache dedicadas.
A fila de tarefas é apoiada por duas tabelas Postgres no Lakebase. A tabela tasks contém uma linha por unidade lógica de trabalho, registrando o status atual da tarefa, informações de concessão, atribuição de agente, extração pai e resultado final. A tabela task_attempts contém uma linha por tentativa de execução, capturando o ID de execução do Databricks Job, o ID de trace do MLflow e os metadados de custo por tentativa. A relação pai-filho oferece suporte a novas tentativas (uma única tarefa pode ter várias tentativas) e preserva a observabilidade no nível da tentativa para atribuição de custos e depuração.
Um par de tabelas Postgres por si só ainda não é uma fila de tarefas. Quatro padrões nativos do Postgres as transformam em uma fila robusta, simultânea, resiliente a falhas e ciente de limites de taxa, adequada para cargas de trabalho agênticas de longa duração.
Uma consulta básica de desenfileiramento pode selecionar a próxima tarefa disponível usando WHERE status = 'enqueued' and LIMIT batch_size. Embora essa consulta identifique corretamente uma tarefa enfileirada, ela não é suficiente quando vários workers estão desenfileirando simultaneamente. Sem o bloqueio de linha, vários workers podem selecionar a mesma tarefa antes que o status seja atualizado.
Adicionar FOR UPDATE SKIP LOCKED torna a remoção da fila (dequeue) segura contra concorrência. Cada worker bloqueia a linha que seleciona, enquanto outros workers pulam essa linha e continuam para a próxima tarefa disponível. Além disso, uma cláusula ORDER BY priority DESC, created_at garante que as tarefas de maior prioridade sejam selecionadas primeiro, preservando a ordenação FIFO dentro de cada nível de prioridade.
A instrução completa, segura contra concorrência e estável em termos de prioridade, é:
Os workers podem ser finalizados no meio de uma tarefa devido à ejeção de VM, condições de falta de memória ou eventos de implantação. Se as tarefas finalizadas continuarem marcadas como em processamento, elas serão retidas indefinidamente. A solução é registrar um lease com expiração no momento da remoção da fila (dequeue):
Um sweeper periódico coloca novamente na fila qualquer tarefa cujo lease_expires_at tenha expirado. As tarefas retidas por workers finalizados são recuperadas automaticamente em poucos minutos, sem a necessidade de um serviço de coordenação externo.
Os endpoints de modelos de visão e LLM geralmente aplicam duas cotas distintas: um limite de requisições por segundo e um limite de tokens por minuto (TPM). Uma única estratégia de controle de vazão (throttling) raramente serve para ambos. O orquestrador suporta três modos, selecionados por agente por meio de configuração.
Limite de concorrência. Um parâmetro MAX_CONCURRENT_TASKS limita o número de tarefas que o orquestrador envia simultaneamente. O limite é aplicado no momento da remoção da fila, contando as linhas atuais em estado PROCESSING na tabela de tarefas:
Se a contagem estiver no limite ou acima dele, nenhuma nova tarefa será removida da fila. Ancorar a verificação na contagem de linhas do banco de dados, em vez de no tamanho da fila do executor local, mantém o limite preciso durante reinicializações de workers, recuperações de lease e implantações com múltiplas réplicas. Este modo é ideal para endpoints limitados por requisições por segundo, onde o uso de tokens de cada tarefa é praticamente uniforme.
Orçamento de tokens. Um parâmetro MAX_TPM limita a taxa de tokens projetada para as tarefas em andamento. O orquestrador estima a contagem de tokens de uma tarefa e soma a taxa de tokens projetada de todas as tarefas em estado PROCESSING. Uma nova tarefa só é removida da fila se a soma mais os tokens projetados da nova tarefa couberem no orçamento.
Limite combinado. Quando ambos MAX_CONCURRENT_TASKS e MAX_TPM estão configurados, o orquestrador aplica a restrição mais rígida. Este modo lida com cargas de trabalho que são limitadas por concorrência sob um regime (muitas tarefas curtas e baratas) e limitadas por tokens sob outro (um único documento muito longo que satura a cota por minuto).
Em todos os três modos, a decisão de controle de vazão é tomada no momento da remoção da fila, dentro da mesma transação que FOR UPDATE SKIP LOCKED. Uma tarefa que não cabe na cota atual permanece na fila e é reconsiderada no próximo ciclo de remoção — sem estado de agendamento separado, sem fila de espera em memória e sem camada de coordenação entre réplicas de workers.
Quando a camada de agentes de IA conclui uma tarefa, ela envia um callback para o orquestrador com o resultado. A entrega de callbacks não é do tipo exatamente-uma-vez (exactly-once): o Databricks pode tentar novamente, as redes podem sofrer interrupções e os proxies podem reenviar. O manipulador de callback foi projetado para ser idempotente: ele aceita os estados PROCESSING e ENQUEUED, e trata tarefas já finalizadas como sem operação (no-ops). Payloads idênticos produzem resultados idênticos, eliminando o risco de cobrança dupla ou processamento duplicado.
Esses quatro padrões combinados resultam em uma fila de tarefas que é correta sob concorrência, durável contra falhas, sensível a limites de taxa sob carga e idempotente sob tentativas de reenvio. A visibilidade em tempo real do sistema em execução é fornecida por um mecanismo separado descrito na próxima seção.
Quando muitos documentos estão em andamento, os operadores precisam de uma visão clara do desempenho do agente, do status da tarefa e do custo da carga de trabalho. Eles não devem precisar fazer consultas contínuas (polling) ao orquestrador ou depender de uma plataforma de métricas separada. O orquestrador integra essa funcionalidade diretamente em um único dashboard exibido pelo mesmo Databricks App que executa o daemon do worker.
O dashboard voltado para o operador apresenta um conjunto de métricas operacionais que, juntas, caracterizam o sistema em execução. Todas métricas suportam filtragem por intervalo de datas, status da tarefa e agente.
As alterações de estado na tabela de tarefas disparam eventos LISTEN/NOTIFY do Postgres. O back-end mantém uma única conexão LISTEN e distribui os eventos via Server-Sent Events (SSE) para os clientes do dashboard conectados. Os navegadores abrem uma conexão EventSource e recebem atualizações em tempo real em aproximadamente um segundo após cada alteração de estado significativa. A implementação não requer Redis, servidor WebSocket ou barramento de mensagens.
O polling é mantido como um fallback permanente em um intervalo padrão de dez segundos. Conexões de streaming por meio de proxies de entrada na nuvem podem descartar bytes silenciosamente sem disparar eventos de erro no lado do cliente; o polling permanente garante que o dashboard permaneça atualizado mesmo nesses casos. Um indicador de UI distingue entre canais live (SSE ativo) e polling (SSE indisponível).
Os dados do dashboard abrangem três fontes com diferentes características de latência: Postgres (instantâneo), API de rastreamento do MLflow (subsegundo) e consultas de warehouse em tabelas de faturamento do sistema (ocasionalmente dezenas de segundos). As consultas rápidas alimentam cada ciclo de atualização; as consultas lentas são executadas apenas sob ação do usuário e retornam de forma otimista com um estado de carregamento até que os resultados estejam disponíveis.
As tabelas de faturamento do sistema Databricks têm escopo de conta: cada job, cada chamada de modelo e qualquer outra aplicação contribuem para as mesmas linhas de system.billing.usage. Sem a definição de escopo, um bloco de "custo de OCR" no nível da aplicação agregaria o uso de todas as chamadas de modelo no workspace.
A solução é registrar quais execuções de Databricks Jobs o orquestrador enviou (rastreadas em tasks.locked_by e task_attempts.run_id) e filtrar a consulta de faturamento para esse conjunto. Um único SQL warehouse pode dar suporte a várias aplicações, e cada dashboard exibe apenas seus próprios gastos.
A mesma arquitetura de consulta se integra naturalmente com filtros definidos pelo operador. Os valores de custo, junto com todas as outras métricas do dashboard, podem ser ainda mais refinados por intervalo de datas, status da tarefa ou agente, respondendo a perguntas como "quanto custaram as tarefas que falharam nos últimos sete dias?" ou "qual foi a mediana de gastos por tarefa para o agente X este mês?" sem sair do dashboard.
Isso torna os valores de custo fáceis de monitorar, alocar e relatar.
Os padrões de Postgres como fila estão bem estabelecidos na comunidade de engenharia de dados. O Lakebase oferece as características operacionais adicionais que tornam o padrão viável como uma arquitetura de produção no Databricks:
Esses recursos eliminam a sobrecarga operacional que normalmente motiva as equipes a adotar message brokers gerenciados em vez de Postgres auto-hospedado para o enfileiramento de tarefas.
Na CLA, o padrão de orquestração descrito aqui oferece suporte a um fluxo de trabalho de processamento de documentos em produção que reduz o tempo de extração de horas para minutos. A arquitetura usa serviços nativos do Databricks com o Lakebase Postgres no centro para gerenciar o enfileiramento, o agendamento e a observabilidade sem a necessidade de sistemas externos. Isso reduz a sobrecarga de integração ao mesmo tempo em que aproveita ao máximo uma plataforma unificada criada para escalar.
Em produção, esse padrão oferece gerenciamento durável de tarefas, controle de prioridade, agendamento com reconhecimento de limite de taxa, visibilidade em tempo real e rastreamento de custos por tarefa. Juntos, esses recursos oferecem uma maneira prática de orquestrar cargas de trabalho de agentes, mantendo a infraestrutura ao redor simples.
Pronto para simplificar a orquestração de agentes de AI em uma única plataforma? Experimente a Databricks Free Edition, crie seu primeiro projeto do Lakebase Postgres e seu primeiro Databricks App e, em seguida, acompanhe a demonstração de 10 minutos do MLflow Tracing para adicionar observabilidade de ponta a ponta ao fluxo de trabalho do seu agente.
(Esta publicação no blog foi traduzida utilizando ferramentas baseadas em inteligência artificial) Publicação original
Assine nosso blog e receba os posts mais recentes diretamente na sua caixa de entrada.