Barramento de Serviço do Azure: filas, tópicos, assinaturas e mensageria confiável para IA
Desacople a entrada de solicitações da latência variável de inferência com nivelamento de carga, consumidores concorrentes, fan-out filtrado, mensagens estruturadas, claim check, liquidação Peek-Lock, idempotência, renovação de bloqueio e recuperação observável da fila de mensagens mortas.
Tempo de estudo sugerido: 125 minutos • Nível intermediário • Reescrita autoral completa com versão resumida de cada tópico, avaliação comentada e laboratório guiado em Python
Por João Ricardo Dutra••Conteúdo autoral completo
1. Cenário de mensageria para IA e objetivos
Imagine uma de análise de documentos cuja inferência demora de poucos segundos a meio minuto. O tráfego chega em rajadas, poucos processadores podem usar o modelo ao mesmo tempo e serviços de notificação, auditoria, métricas e qualidade precisam do resultado. Uma cadeia síncrona prolonga a espera, exige disponibilidade simultânea e permite que uma entrada lenta ou inválida ocupe capacidade útil.
O do insere um broker durável entre esses componentes. A confirma a aceitação rapidamente, os workers drenam o trabalho em um ritmo controlado e os consumidores posteriores recebem cópias independentes. Uma entrada problemática permanece diagnosticável.
Aplicar nivelamento de carga, consumidores concorrentes, desacoplamento temporal e publicação/assinatura em IA.
Escolher filas ou tópicos com assinaturas, sessões e filtros no broker.
Projetar corpos , propriedades, correlação, lotes, expiração e check.
Processar com Peek-Lock, liquidação explícita, idempotência, renovação de bloqueio e recuperação pela DLQ.
Testar o fluxo em Python com -servicebus e .
Resumo do tópico
O broker separa a aceitação rápida da solicitação do processamento variável de IA e dá um caminho operacional explícito a entregas, repetições e falhas.
2. Namespaces, entidades, protocolos e identidade
O é um broker de mensagens empresarial totalmente gerenciado. O é o limite administrativo e de rede que contém filas, tópicos e assinaturas e expõe um como <>.servicebus.windows.net. AMQP 1.0 é o protocolo principal dos SDKs modernos e viabiliza liquidação, transações, controles de ordem e detecção de duplicidades.
É possível autenticar por SAS ou . Em cargas no , prefira e a função de dados mais restrita: Data Sender, Data Receiver ou . Isso elimina credenciais embutidas; depois da migração dos clientes, a autenticação local pode ser desativada.
Blocos fundamentais
Elemento
Responsabilidade
Define entidades, , camada, capacidade, rede, diagnóstico e política de autenticação
Fila
Armazena trabalho duravelmente para processamento ponto a ponto ou concorrente
Tópico
Recebe uma publicação e distribui cópias compatíveis às assinaturas
Assinatura
Funciona como fila virtual independente e pode ter regras e filtros
Cliente AMQP 1.0
Envia, recebe, bloqueia e liquida mensagens pelo
Resumo do tópico
Trate o como limite de segurança e capacidade, escolha as entidades corretas e use com privilégio mínimo.
3. Desacople a entrada da solicitação da inferência
Em uma solicitação-resposta assíncrona, a valida o envelope, grava o estado, envia a mensagem e devolve um identificador da operação. Depois, um worker recebe, executa a inferência, persiste o resultado e publica um evento de conclusão ou permite consulta pelo cliente. Produtor e consumidor deixam de depender do mesmo instante de implantação e escala.
A arquitetura não elimina a latência: ela a torna visível e controlável. O produto deve representar os estados aceito, em execução, concluído, falho e expirado, proteger a consulta de estado e definir e notificação. A fila não substitui o banco de resultados.
Resumo do tópico
A mensageria assíncrona transforma uma inferência longa em operação rastreável, com estado e resultado persistidos e evolução independente dos componentes.
4. Nivele a carga com uma fila durável
O nivelamento absorve uma rajada curta na fila enquanto um conjunto estável de workers consome na vazão sustentável. Isso protege GPUs, memória e de modelo limitados e evita manter computação suficiente para um pico raro.
A fila exige capacidade e monitoramento. Um backlog crescente eleva o tempo de conclusão e pode perder valor de negócio. Defina capacidade, , alertas e limites de escala com base nas taxas de chegada e processamento, idade aceitável e comportamento de falha. Torne a contrapressão perceptível aos chamadores.
Diagrama autoral: a fila durável suaviza rajadas da e distribui mensagens entre workers de IA com escala independente.
Resumo do tópico
Use a fila como controlado entre demanda irregular e capacidade finita, com limites explícitos de tamanho e idade do backlog.
5. Escale com consumidores concorrentes e desacoplamento temporal
Várias instâncias recebem da mesma fila. O bloqueia cada mensagem para um receptor por vez, distribuindo o trabalho sem um despachante central. , ou podem adicionar instâncias conforme o backlog cresce.
Se um worker falhar antes da liquidação, o bloqueio expira e a mensagem volta a ficar disponível. Essa recuperação também permite duplicidades, por isso os efeitos posteriores devem ser idempotentes. O armazenamento durável também mantém trabalho durante uma implantação ou interrupção breve.
Resumo do tópico
Consumidores concorrentes distribuem trabalho horizontalmente; armazenamento durável desacopla disponibilidade e idempotência torna a reentrega segura.
6. Use profundidade da fila como contrapressão e sinal de escala
Mensagens ativas, idade da mais antiga, taxas de entrada e conclusão, latência, falhas e mensagens mortas descrevem melhor a saúde do que CPU isolada. Crescimento contínuo indica que a chegada supera a conclusão; fila sempre vazia pode significar baixa latência ou excesso de capacidade.
O coleta métricas e dispara alertas. KEDA nos ou no e gatilhos do convertem backlog em réplicas. Configure mínimo, máximo, cooldown, concorrência e limites do de modelo em conjunto para não apenas deslocar o gargalo.
Resumo do tópico
Dimensione com backlog, idade, vazão e limites posteriores em conjunto; profundidade é sinal de pressão, não um plano completo.
7. Escolha Standard ou Premium e respeite limites
Decisões de camada
Questão
Standard
Premium
Capacidade
Infraestrutura compartilhada
Unidades de mensageria dedicadas e maior isolamento
Filas, tópicos e assinaturas
Com suporte
Com suporte
Mensagem única
256 KB
1 MB por padrão; até 100 MB por entidade via AMQP quando configurado
/SBMP
Dentro do limite da camada
Até 1 MB por mensagem
256 KB
Até 1 MB, mesmo com mensagens grandes habilitadas
Rede e resiliência
Controles essenciais
privados, opções de rede virtual e zonas quando disponíveis
A cota inclui corpo e propriedades e varia com camada, protocolo e configuração; confira a documentação vigente. Documentos, imagens, áudio e artefatos de modelo normalmente pertencem ao com check, mesmo quando uma mensagem grande seria tecnicamente aceita.
Resumo do tópico
Escolha a camada por isolamento, rede, disponibilidade e vazão medida; use as cotas como limites de validação, não como meta de tamanho.
8. Escolha fila ou tópico com assinaturas
Seleção da entidade
Necessidade
Use
Motivo
Um worker deve processar cada solicitação
Fila
Consumidores concorrentes compartilham o trabalho
Vários serviços precisam do mesmo resultado
Tópico + assinaturas
Cada assinatura compatível recebe uma cópia
Um produtor e uma responsabilidade
Fila
Ciclo de vida mais simples
Notificação, auditoria, métricas e qualidade
Tópico + assinaturas
Escala e falha independentes
Adicionar consumidores sem mudar o publicador
Tópico + assinaturas
Publicador usa um tópico estável
O receptor não lê o tópico diretamente: lê uma assinatura, que funciona como fila virtual e também pode ter consumidores concorrentes. Não use três consumidores na mesma fila quando os três precisam ver a mensagem; somente um deles a receberá.
Resumo do tópico
Use fila para uma responsabilidade de processamento e tópico com assinaturas quando várias responsabilidades precisarem de cópias independentes.
9. Preserve a ordem por fluxo com sessões
A ordem de chegada não garante a conclusão ordenada entre workers. Sessões agrupam mensagens por session_id e concedem um bloqueio exclusivo da sessão, oferecendo FIFO dentro do grupo e paralelismo entre grupos. Extração, classificação e resumo podem compartilhar o identificador do documento.
A entidade deve nascer com sessões habilitadas e toda mensagem precisa de session_id. Sessões reduzem a concorrência de chaves quentes e acrescentam estado e bloqueios; use-as somente quando a ordem for requisito de correção. Dependências complexas talvez sejam mais claras em um orquestrador.
Resumo do tópico
Sessões fornecem processamento exclusivo e ordenado por chave de negócio, mantendo paralelismo entre chaves.
10. Distribua resultados com filtros de assinatura
Cada assinatura começa com uma TrueFilter que aceita tudo. Troque-a ou complemente-a com filtros SQL sobre propriedades de sistema e aplicação; filtros de correlação são eficientes para correspondências exatas. FalseFilter não aceita mensagens e ajuda quando todas as regras serão explícitas.
Propriedades como priority, model_name, document_type, tenant e review_required permitem roteamento sem interpretar . Estabilize nomes e tipos, teste regras sobrepostas e remova a regra verdadeira padrão quando ela inviabilizar a seleção. Filtro controla entrega, não autorização.
Diagrama autoral: uma publicação gera cópias independentes para notificação, auditoria, métricas e qualidade segundo regras do broker.
Resumo do tópico
Roteie publicações com propriedades estáveis e regras testadas; cada assinatura mantém backlog, repetição, escala e falha próprios.
11. Gerencie remetentes e receptores no do Python
O pacote -servicebus oferece ServiceBusClient, remetentes de fila/tópico e receptores de fila/assinatura. Gerenciadores de contexto fecham links AMQP de forma previsível. Reutilize clientes e links duradouros quando adequado e siga a referência atual do , pois e versões de Python evoluem.
from azure.identity import DefaultAzureCredential
from azure.servicebus import ServiceBusClient, ServiceBusMessage
namespace = "<namespace>.servicebus.windows.net"
credential = DefaultAzureCredential()
with ServiceBusClient(namespace, credential) as client:
with client.get_queue_sender("inference-requests") as sender:
sender.send_messages(ServiceBusMessage(
'{"request_id":"req-917","model":"document-analyzer"}',
content_type="application/json",
message_id="req-917",
correlation_id="trace-6d13",
application_properties={"priority": "high"},
))
DefaultAzureCredential funciona com a identidade do desenvolvedor e com sem alterar a lógica. Separe as funções Sender e Receiver. Use get_topic_sender() para tópicos e get_subscription_receiver(topic_name, subscription_name) para assinaturas.
Resumo do tópico
Trate clientes do como recursos gerenciados, reutilize conexões e autentique cada componente com identidade restrita.
12. Estruture corpos e propriedades de mensagens de IA
A mensagem contém corpo, propriedades de aplicação e propriedades do sistema. é adequado para request_id, modelo, temperature, max_tokens e referências à entrada ou ao contexto. Defina content_type como application/ e versione o contrato para que incompatibilidades sejam rejeitadas ou transformadas conscientemente.
Lugar correto para cada dado
Local
Exemplos
Finalidade
Corpo
Referência ao prompt/documento, parâmetros, entrada do fluxo
Valide esquema, faixas, modelos permitidos, e autorização antes da inferência cara. Não coloque segredos no corpo ou propriedades; ambos entram na cota e podem ser inspecionados por operadores autorizados.
Resumo do tópico
Mantenha o contrato versionado no corpo, valores leves de roteamento nas propriedades da aplicação e semântica de entrega nas propriedades do sistema.
13. Correlacione, rastreie, deduplique e processe com idempotência
message_id identifica a mensagem e alimenta a detecção de duplicidade dentro da janela configurada. correlation_id conecta , fila, inferência, publicação do resultado e . Para OpenTelemetry, propague traceparent e tracestate nas propriedades de aplicação.
A detecção de duplicidade protege reenvios com o mesmo message_id, mas não impede reentrega depois de falha do receptor ou liquidação ambígua. Registre um identificador de negócio em armazenamento durável e torne gravações, notificações e cobranças repetíveis sem efeito extra. Peek-Lock continua sendo entrega pelo menos uma vez; idempotência produz um efeito de negócio efetivamente único quando bem projetada.
Resumo do tópico
Use message_id contra reenvios, correlação e contexto de para observabilidade e idempotência durável contra reentregas.
14. Aplique check a cargas grandes
Um documento de 500 MB não cabe em uma mensagem, e cargas grandes permitidas ainda reduzem a vazão. Envie o documento ao privado e publique apenas uma referência opaca, , tamanho, tipo de mídia, modelo e identificadores. O worker autorizado recupera e valida o objeto.
# 1. Upload the large document to private Azure Blob Storage.
blob_uri = upload_with_managed_identity(document_bytes)
# 2. Send only the claim check and routing metadata.
message = ServiceBusMessage(
json.dumps({
"request_id": request_id,
"blob_uri": blob_uri,
"sha256": payload_hash,
"model": "document-analyzer",
}),
content_type="application/json",
message_id=request_id,
correlation_id=correlation_id,
)
sender.send_messages(message)
# 3. The authorized consumer retrieves, validates, and processes the blob.
Prefira e rede privada; se usar SAS, restrinja escopo e duração. Defina propriedade, retenção, repetição e exclusão para não deixar blobs órfãos nem apagar uma entrada antes de um . check também evita expor conteúdo sensível a intermediários.
Resumo do tópico
Guarde entradas grandes ou sensíveis em armazenamento de objetos protegido e envie uma referência pequena e verificável com ciclo de vida coordenado.
15. Controle validade com e vazão com lotes
time_to_live representa por quanto tempo o trabalho continua útil. Recomendações em tempo real podem expirar cedo; análises em lote, mais tarde. Mensagens expiradas podem ir à fila de mensagens mortas quando a opção está habilitada. Mensagens adiadas têm comportamento especial e não devem virar armazenamento permanente.
ServiceBusMessageBatch agrupa mensagens até o limite calculado. Quando add_message não comporta a próxima, envie o lote, crie outro e tente novamente. Uma mensagem individual grande ainda exige check. O lote reduz viagens de rede, mas não cria uma transação de negócio única.
from azure.servicebus.exceptions import MessageSizeExceededError
batch = sender.create_message_batch()
for payload in payloads:
message = ServiceBusMessage(json.dumps(payload))
try:
batch.add_message(message)
except MessageSizeExceededError:
sender.send_messages(batch)
batch = sender.create_message_batch()
batch.add_message(message)
if len(batch) > 0:
sender.send_messages(batch)
Resumo do tópico
Defina pelo valor de negócio e agrupe mensagens pequenas até o limite do , sem confundir lote de transporte com transação.
16. Receba com Peek-Lock e liquide conscientemente
Recebimento e liquidação
Escolha
Efeito
Uso
Receive-and-
Remove na entrega; falha pode perder trabalho
Telemetria não crítica
Peek-Lock
Bloqueia e remove apenas após complete
Padrão para inferência importante
Complete
Sucesso e remoção definitiva
Depois dos efeitos duráveis
Abandon
Libera para nova tentativa e incrementa entregas
Falha transitória
Dead-letter
Move para DLQ com diagnóstico
Entrada permanentemente inválida
Defer
Mantém, mas exige sequence_number
Dependência conhecida ou ordem intencional
from azure.servicebus import ServiceBusReceiveMode
with client.get_queue_receiver(
queue_name="inference-requests",
receive_mode=ServiceBusReceiveMode.PEEK_LOCK,
max_wait_time=30,
) as receiver:
for message in receiver:
try:
payload = json.loads(str(message))
validate(payload)
process_idempotently(payload, str(message.message_id))
receiver.complete_message(message)
except PermanentPayloadError as error:
receiver.dead_letter_message(
message,
reason="InvalidPayload",
error_description=str(error),
)
except TransientDependencyError:
receiver.abandon_message(message)
Liquide somente após concluir efeitos duráveis. durante a liquidação é ambíguo: o broker pode ter aplicado a operação sem a confirmação chegar ao cliente, reforçando a necessidade de idempotência.
Resumo do tópico
Peek-Lock prioriza recuperação: complete o sucesso, abandone falhas transitórias, envie falhas permanentes à DLQ e adie apenas com plano de recuperação.
17. Renove bloqueios e opere a fila de mensagens mortas
O bloqueio padrão da entidade dura um minuto e pode chegar a cinco. Para processamento legítimo mais longo, renove-o manualmente ou use AutoLockRenewer por período limitado. Evite receber ou pré-buscar mais mensagens do que o worker concluirá. Para tarefas sempre muito longas, registre o trabalho, conclua a mensagem rapidamente e use uma máquina de estados separada.
from azure.servicebus import AutoLockRenewer
with AutoLockRenewer() as renewer:
with client.get_queue_receiver("inference-requests") as receiver:
for message in receiver:
renewer.register(
receiver,
message,
max_lock_renewal_duration=600,
)
run_long_inference(message)
receiver.complete_message(message)
Cada fila e assinatura tem uma DLQ. Ao ultrapassar maxDeliveryCount — 10 por padrão — a mensagem chega com MaxDeliveryCountExceeded; a aplicação também pode registrar motivo e descrição próprios. A DLQ retém tudo até liquidação explícita. Alerte por quantidade e idade, corrija a causa e só então reproduza por processo idempotente aprovado.
from azure.servicebus import ServiceBusSubQueue
with client.get_queue_receiver(
queue_name="inference-requests",
sub_queue=ServiceBusSubQueue.DEAD_LETTER,
max_wait_time=10,
) as dlq_receiver:
for message in dlq_receiver:
inspect(
message.dead_letter_reason,
message.dead_letter_error_description,
message.delivery_count,
message.correlation_id,
)
# Re-submit only after fixing the cause and preserving idempotency.
replay_if_approved(message)
dlq_receiver.complete_message(message)
Diagrama autoral: uma mensagem bloqueada é concluída, abandonada, adiada ou enviada à DLQ; a reprodução ocorre apenas depois do reparo.
Resumo do tópico
Renove bloqueios para processamento longo limitado, controle tentativas e trate a DLQ como fluxo observável de reparo.
18. Laboratório, revisão da avaliação e checklist de produção
O exercício de origem cria e aplicativo Flask em Python, usa fila com Peek-Lock, inspeciona uma entrada inválida na DLQ e distribui resultados a assinaturas filtradas. Reserve cerca de 30 minutos, uma assinatura do ,, Python 3.12 ou posterior e atual.
az group create --name ai200-servicebus-rg --location eastus
az servicebus namespace create --resource-group ai200-servicebus-rg --name <globally-unique-namespace> --location eastus --sku Standard
az servicebus queue create --resource-group ai200-servicebus-rg --namespace-name <namespace> --name inference-requests --max-delivery-count 5
az servicebus topic create --resource-group ai200-servicebus-rg --namespace-name <namespace> --name inference-results
python -m venv .venv
python -m pip install --upgrade azure-identity azure-servicebus flask
Crie , fila, tópico e assinaturas; configure tentativas e filtros explicitamente.
Atribua funções Sender ou Receiver às identidades gerenciadas e use DefaultAzureCredential.
Envie válido e uma mensagem inválida; processe o trabalho válido idempotentemente com Peek-Lock.
Confirme a DLQ, corrija a causa e reproduza a mensagem uma única vez.
Publique um resultado com prioridade e prove que somente assinaturas compatíveis o recebem.
Observe backlog, idade, conclusões, repetições, perda de bloqueio e DLQ; exclua os recursos ao terminar.
Respostas da avaliação
Pergunta
Resposta
Motivo
Três serviços precisam de cada resultado
Tópico com três assinaturas
Cada assinatura ganha uma cópia
Falha do worker não pode perder a solicitação
Peek-Lock
Trabalho não liquidado volta a ficar disponível
Décima falha alcança o limite
DLQ com MaxDeliveryCountExceeded
O broker isola a mensagem problemática
Documento de 500 MB
check no
O broker transporta somente a referência
Finalidade de correlation_id
Rastreamento ponta a ponta
Liga etapas, e resultado
Antes da produção, confirme camada e cotas, infraestrutura como código, privilégio mínimo, autenticação local, rede privada, validação de contrato, , janela de duplicidade, sessões, filtros, armazenamento de idempotência, bloqueio, prefetch, , responsável pela DLQ, alertas, runbooks, custo e testes de carga no real.
O projeto completo combina entidade correta, mensagens pequenas e versionadas, identidade, filtros, correlação, Peek-Lock, idempotência, controle de bloqueio, recuperação observável da DLQ e teste ponta a ponta.