Redis Gerenciado pelo Azure: Pub/Sub, Streams, eventos e filas de tarefas
Difunda eventos de IA ao vivo, coordene trabalho durável com Redis Streams, recupere entregas abandonadas, escale grupos de consumidores e combine os dois modelos com segurança.
Tempo de estudo sugerido: 110 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. Desacoplar um de IA em tempo real
Imagine uma plataforma jurídica que executa OCR, reconhecimento de entidades, classificação e geração de embeddings para centenas de documentos simultâneos. Uma síncrona manteria cada requisição aberta por vários segundos; e tentativas criados à mão seriam frágeis. A mensageria permite responder rápido ao upload, escalar cada etapa separadamente e enviar o progresso a painéis ao vivo.
O e o oferecem Pub/Sub e Redis . Pub/Sub distribui uma notificação transitória aos ouvintes conectados. conserva o trabalho para consumidores coordenados processarem e confirmarem. Os exemplos usam Python e redis-py; como a da biblioteca muda, confira as assinaturas na versão instalada.
Difundir um evento a vários serviços de IA com canais e padrões.
Criar uma fila durável com e grupos de consumidores.
Recuperar explicitamente trabalho pendente após falha.
Escolher difusão, distribuição coordenada ou arquitetura híbrida.
Executar um exercício Flask com e CLI do .
Uma mesma entrada pode produzir uma notificação ao vivo e uma unidade de trabalho durável.
Resumo do tópico
Separe a requisição rápida do processamento assíncrono e use Pub/Sub para difusão ao vivo e para trabalho retido.
2. Distinguir evento, mensagem e tarefa
A intenção determina o modelo.
Carga
Significado
Expectativa
Evento
Fato que já ocorreu, como model_updated
Zero, uma ou várias reações independentes
Notificação
Atualização breve, como prediction_ready
Somente ouvintes conectados podem precisar dela
Tarefa ou comando
Pedido de trabalho, como analyze_document
Um trabalhador precisa concluir, repetir ou registrar falha
O publicador não deve conhecer cada receptor. O intermediário desacopla implantação e escala; ID de correlação, tipo, versão do esquema, horário e uma referência para cargas grandes preservam rastreabilidade. Mantenha documentos e artefatos de modelo fora da mensagem e transporte identificadores.
Resumo do tópico
Eventos anunciam fatos; tarefas solicitam trabalho. Modele a carga e a exigência de entrega antes de escolher o recurso Redis.
3. Entender a difusão por canais Redis Pub/Sub
O publicador executa PUBLISH com canal e carga. O Redis envia imediatamente uma cópia a cada assinante ativo, sem que o emissor conheça os ouvintes. Assim, serviços de sentimento, intenção e contexto reagem em paralelo a um evento de nova conversa.
O desacoplamento oferece baixa latência e alta vazão em muitos canais. Ele serve para coordenação em tempo real, invalidação de , avisos de atualização de modelos ou embeddings, status de treinamento, previsões e telemetria útil apenas enquanto o consumidor está online.
Pub/Sub é um barramento de difusão para ouvintes conectados, não uma fila durável.
Resumo do tópico
PUBLISH cria difusão um-para-muitos: cada assinante conectado recebe o evento de forma independente.
4. Projetar canais e assinaturas por padrão
Canais são strings. Um previsível mostra domínio, assunto e evento, como ai:models:updated. Acrescente locatário ou entidade apenas para isolamento ou entrega direcionada; canais ilimitados por usuário ficam difíceis de operar.
SUBSCRIBE nomeia canais exatos. PSUBSCRIBE aceita padrões glob e entrega uma pmessage com padrão, canal real e carga. Padrões amplos simplificam descoberta, mas aumentam tráfego e podem expor eventos indevidos; autorize e documente os namespaces.
subscription = client.pubsub()
subscription.psubscribe("ai:*")
for event in subscription.listen():
if event["type"] == "pmessage":
handle_ai_event(
pattern=event["pattern"],
channel=event["channel"],
payload=event["data"],
)
Resumo do tópico
Use namespaces estáveis, assinaturas exatas para tráfego restrito e padrões limitados para famílias relacionadas.
5. Considerar entrega no máximo uma vez e pressão de fluxo
Sem persistência: a mensagem não fica disponível para repetição.
Entrega no máximo uma vez: assinante desconectado, reiniciando ou com falha perde o evento.
Sem confirmação: a contagem retornada não prova que o processamento de negócio terminou.
Sem contrapressão de fila: o assinante lento precisa enfileirar, descartar ou desconectar no próprio processo.
Sem histórico: ordenação, auditoria e exigem outro armazenamento ou modelo.
Essas propriedades ajudam quando o evento é substituível ou efêmero. Elas são inadequadas se perder um documento ou cobrar duas vezes for relevante. Não transforme Pub/Sub em fila atribuindo a mesma assinatura a todos os trabalhadores: cada um receberá a tarefa e multiplicará o custo.
Resumo do tópico
Pub/Sub prioriza velocidade e fan-out em vez de durabilidade, confirmação, e divisão coordenada da carga.
6. Aplicar Pub/Sub aos cenários corretos de IA
Cenários de difusão.
Cenário
Por que fan-out serve
Invalidação de modelo ou embedding
Cada instância limpa seus dados locais obsoletos
Configuração e sinalizadores
Todos os serviços carregam a mesma alteração
Previsão pronta
, painéis e telemetria reagem separadamente
Métricas de IA
Alertas, painéis e consomem o mesmo sinal de maneiras distintas
Ações heterogêneas
Análise, cobrança e recomendações executam ações diferentes após uma interação
A falha de um assinante não deve bloquear publicador nem outros consumidores. Se o fato também exigir uma ação garantida, persista a tarefa separadamente — por exemplo, em um — e deixe Pub/Sub apenas para a visão ao vivo.
Resumo do tópico
Escolha Pub/Sub quando todos os serviços conectados devem ver o mesmo evento substituível e reagir de forma independente.
7. Publicar e ouvir com redis-py
A conexão considera a porta 10000 do e um provedor já configurado. Em produção, rode o listener bloqueante em trabalhador ou thread dedicada, trate reconexões, valide esquemas e encerre a assinatura corretamente.
import redis
client = redis.Redis(
host="<cache-name>.<region>.redis.azure.net",
port=10000,
ssl=True,
credential_provider=entra_provider,
decode_responses=True,
)
client.publish("ai:models:updated", "summarizer:v3")
subscription = client.pubsub()
subscription.subscribe("ai:models:updated", "ai:embeddings:refresh")
for event in subscription.listen():
if event["type"] == "message":
handle_event(event["channel"], event["data"])
Verifique event["type"], pois quadros de confirmação de assinatura também chegam. decode_responses=True converte canais e cargas em strings; desative-o para binários. O inteiro de publish() informa quantos clientes receberam o quadro, não quantos concluíram a ação.
Resumo do tópico
Use publish(), pubsub(), subscribe(), psubscribe() e listen() com quadros explícitos, conexão segura e ciclo resiliente.
8. Encaminhar eventos Redis ao navegador
O navegador normalmente não acessa Redis diretamente. Um serviço FastAPI ou Flask confiável assina em segundo plano e encaminha eventos autorizados por . Ele deve mapear usuários a tópicos permitidos, limitar filas por , fechar conexões ociosas e impedir vazamento entre locatários.
from fastapi import FastAPI, WebSocket
import redis.asyncio as redis
app = FastAPI()
async def forward_predictions(socket: WebSocket):
client = redis.Redis(
host="<cache-name>.<region>.redis.azure.net",
port=10000,
ssl=True,
credential_provider=entra_provider,
decode_responses=True,
)
async with client.pubsub() as subscription:
await subscription.subscribe("ai:predictions:ready")
async for event in subscription.listen():
if event["type"] == "message":
await socket.send_json({"type": "prediction", "data": event["data"]})
Uma assinatura Redis por navegador pode esgotar conexões. Um de produção costuma compartilhar uma assinatura por processo, demultiplexar mensagens aos locais e publicar somente status compacto; o resultado durável permanece atrás de autenticadas.
Resumo do tópico
O converte fan-out Redis em atualizações do navegador com autorização, e limites.
9. Modelar trabalho durável como Redis
Um Redis é uma sequência append-only de campos e valores. XADD acrescenta a tarefa e retorna um ID ordenado no tempo, como 1699980000000-0. Diferente de Pub/Sub, a entrada fica disponível até ser removida ou aparada. A de upload enfileira a inferência e devolve o ID enquanto os trabalhadores processam.
task_id = client.xadd(
"ai:inference:queue",
{
"user_id": "42",
"model": "summarizer",
"prompt": "Summarize document 917",
"priority": "high",
},
maxlen=10_000,
approximate=True,
)
# Return immediately from the API instead of waiting for inference.
return {"task_id": task_id, "status": "queued"}
servem para inferência, extrair-analisar-resumir-incorporar, operações longas, histórico e recuperação. O ID é uma boa chave de correlação e idempotência, mas não torna o efeito de negócio exatamente uma vez.
Resumo do tópico
XADD cria um registro ordenado e retido, permitindo que a responda antes do processamento de IA.
10. Distribuir tarefas com grupos de consumidores
XGROUP CREATE estabelece o grupo e MKSTREAM pode criar o . Trabalhadores compartilham o nome do grupo, mas usam nomes de consumidor exclusivos e sensíveis a maiúsculas. XREADGROUP com > pede entradas ainda não entregues naquele grupo. O Redis reparte novas entradas entre leitores ativos, permitindo escala horizontal sem balanceador no código.
import os
import redis
STREAM = "ai:inference:queue"
GROUP = "inference-workers"
try:
client.xgroup_create(STREAM, GROUP, id="0", mkstream=True)
except redis.ResponseError as error:
if "BUSYGROUP" not in str(error):
raise
consumer = f"{os.getenv('HOSTNAME', 'local')}-{os.getpid()}"
while True:
batches = client.xreadgroup(
groupname=GROUP,
consumername=consumer,
streams={STREAM: ">"},
count=5,
block=5000,
)
for _, tasks in batches:
for task_id, fields in tasks:
process_idempotently(task_id, fields)
client.xack(STREAM, GROUP, task_id)
XACK remove a entrega concluída da Pending Entries List (PEL) do grupo; não apaga necessariamente a entrada. Ler com ID como 0 acessa pendências daquele consumidor, não trabalho novo. Vários grupos podem consumir o mesmo de forma independente.
Resumo do tópico
Grupos dividem novas entradas entre trabalhadores únicos e XACK registra o processamento bem-sucedido.
11. Recuperar entregas com falha explicitamente
Se um trabalhador morre depois da entrega e antes de XACK, o Redis mantém a entrada na PEL. Ele não a reatribui nem repete automaticamente. O aplicativo precisa consultar XPENDING e reivindicar itens ociosos com XCLAIM ou XAUTOCLAIM. Um consumidor reiniciado também pode ler sua própria pendência antes de pedir novidades.
# Observe entries that were delivered but not acknowledged.
pending = client.xpending_range(
"ai:inference:queue", "inference-workers", "-", "+", count=100
)
# Redis keeps abandoned entries pending; your recovery loop must claim them.
next_id, claimed, deleted = client.xautoclaim(
"ai:inference:queue",
"inference-workers",
"recovery-worker-1",
min_idle_time=300_000,
start_id="0-0",
count=25,
)
for task_id, fields in claimed:
process_idempotently(task_id, fields)
client.xack("ai:inference:queue", "inference-workers", task_id)
A recuperação produz comportamento pelo menos uma vez: um pode fazer outro trabalhador reivindicar o item enquanto o primeiro termina. Torne efeitos idempotentes pelo ID do ou por chave de negócio, limite tentativas e mova mensagens problemáticas para um de dead letter com diagnóstico.
Resumo do tópico
A PEL preserva trabalho não confirmado, mas o aplicativo deve reivindicar itens ociosos e tolerar nova entrega.
12. Monitorar e limitar a retenção do
XINFO mostra comprimento e IDs; XINFO GROUPS apresenta atraso e pendências; XINFO CONSUMERS revela consumidores, ociosidade e propriedade. Alerte para atraso crescente, pendências antigas, muitas tentativas, troca de trabalhadores, memória, latência e erros.
Apare por comprimento máximo com XADD MAXLEN ou XTRIM; o modo aproximado custa menos.
Mantenha cargas pequenas e documentos ou embeddings em armazenamento durável.
Defina retenção suficiente para auditoria e recuperação; ilimitado não expira sozinho.
Persistência e confirmação adicionam latência e código em relação a Pub/Sub.
Teste encerramento, reconexão, mensagens problemáticas e recuperação sob carga.
Resumo do tópico
Observe , grupo e consumidores e limite entradas para que confiabilidade não gere crescimento sem controle.
13. Escolher difusão ou distribuição coordenada
Guia de decisão.
Requisito
Pub/Sub
com grupo
Quem recebe?
Cada assinante conectado
Um consumidor de cada grupo
Receptor offline
Perde a mensagem
Lê trabalho retido ou pendente depois
Confirmação e repetição
Não integradas
XACK e reivindicação feita pelo aplicativo
Escala de trabalhadores
Duplica trabalho
Divide entradas
e histórico
Indisponíveis
Disponíveis durante a retenção
Latência e complexidade
Menores
Um pouco maiores
Use Pub/Sub para invalidação, configuração, métricas e notificações descartáveis. Use para jobs, inferência e recuperáveis. Se transações, roteamento avançado, dead letter ou garantias entre sistemas dominarem, compare um serviço de mensageria especializado.
A intenção de entrega, não apenas a vazão, define o mecanismo.
Resumo do tópico
Difunda fatos com Pub/Sub, coordene trabalho durável com e avalie um broker dedicado para requisitos mais amplos.
14. Combinar e Pub/Sub na mesma arquitetura
Muitos sistemas precisam dos dois. Ao receber um documento, grave a tarefa durável no e publique um evento efêmero de recebimento. Um membro do grupo processa cada entrega, enquanto , análise e monitoramento veem as difusões. Etapas seguintes podem escrever em novos e publicar progresso.
import json
# Durable work: one worker in the group handles each delivery.
task_id = client.xadd("ai:documents:queue", {
"document_id": "doc-917",
"requested_by": "user-42",
})
# Ephemeral fan-out: every connected observer sees the status event.
client.publish("ai:documents:events", json.dumps({
"event": "document_received",
"task_id": task_id,
}))
As duas gravações não formam uma transação de negócio por padrão. Se perder a notificação for aceitável, trate-a como melhor esforço. Caso precisem de consistência, use outbox ou outro registro durável. Repita o mesmo ID de correlação em , Pub/Sub, e resultados.
Resumo do tópico
O modelo híbrido usa como job oficial e Pub/Sub como observabilidade efêmera, unidos por correlação e regras de consistência.
15. Laboratório guiado: publicar e assinar no
O exercício de origem estima 30–40 minutos e cria uma página Flask Python para publicar e assinar em tempo real. Use recurso descartável, não confirme credenciais no código e remova o ambiente ao final.
Prepare assinatura do ,, Python 3.12 ou posterior, a CLI do mais recente e a extensão redisenterprise instalada com az extension add --name redisenterprise.
Baixe o projeto inicial, crie ambiente virtual isolado e instale dependências fixadas.
Crie o e conceda à identidade de desenvolvimento o acesso necessário.
Conecte por com e use a cadeia normal de credenciais do .
Complete funções para publicar evento, publicar em todos os canais configurados e formatar quadros recebidos.
Execute o listener em thread de segundo plano para manter o Flask responsivo.
Assine canais exatos e padrões, rode a página e confirme mensagens ao vivo.
Teste desconexão temporária para observar a ausência de e exclua o recurso.
Resumo do tópico
O laboratório valida provisionamento seguro, listener em segundo plano, assinaturas exatas e por padrão, semântica e limpeza.
16. Revisão da avaliação e checklist de produção
Pub/Sub atende assinantes ativos; retém entradas para consumo posterior ou coordenado.
Use e grupos quando o precisa de confirmação e repetição gerenciada pelo aplicativo.
XADD acrescenta uma entrada ao .
Pub/Sub serve melhor à difusão de status para clientes conectados.
XREADGROUP com > divide novas entradas no grupo; idempotência protege contra nova entrega.
Versione esquemas e inclua IDs de correlação e idempotência.
Autorize namespaces e deixe cargas grandes fora da mensagem.
Dê nome único a cada consumidor.
Confirme somente depois do efeito de negócio.
Monitore idade da PEL, atraso, repetição, memória, latência e desconexões.
Implemente XAUTOCLAIM ou XPENDING com XCLAIM; Redis não recupera sozinho.
Apare e trate tarefas problemáticas com estratégia de dead letter.