Redis Gerenciado pelo Azure: Pub/Sub, Streams, eventos e filas de tarefas
Voltar para a trilha AI-200
AI-200Capítulo 16

Estudo para a Certificação Microsoft AI-200

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

Escudo neon Microsoft Certified AI-200 com canais Pub/Sub, Streams, grupos de consumidores, filas de tarefas e recuperação no Redis Gerenciado pelo Azure

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 API de IA envia eventos transitórios por Redis Pub/Sub e tarefas duráveis por Redis Streams a observadores independentes e trabalhadores coordenados.
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.
CargaSignificadoExpectativa
EventoFato que já ocorreu, como model_updatedZero, uma ou várias reações independentes
NotificaçãoAtualização breve, como prediction_readySomente ouvintes conectados podem precisar dela
Tarefa ou comandoPedido de trabalho, como analyze_documentUm 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.

Um publicador envia a um canal Redis Pub/Sub que copia o evento para quatro assinantes conectados; um assinante offline não o recebe.
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.

ai:conversations:42
ml:model:predictions
ai:training:status
embeddings:cache:refresh

# Pattern subscriptions
ai:conversations:*
ml:models:*:status
ai:*:predictions

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árioPor que fan-out serve
Invalidação de modelo ou embeddingCada instância limpa seus dados locais obsoletos
Configuração e sinalizadoresTodos os serviços carregam a mesma alteração
Previsão pronta, painéis e telemetria reagem separadamente
Métricas de IAAlertas, painéis e consomem o mesmo sinal de maneiras distintas
Ações heterogêneasAná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.

stream = client.xinfo_stream("ai:inference:queue")
groups = client.xinfo_groups("ai:inference:queue")
consumers = client.xinfo_consumers("ai:inference:queue", "inference-workers")

client.xtrim(
    "ai:inference:queue",
    maxlen=10_000,
    approximate=True,
)
  • 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.
RequisitoPub/Sub com grupo
Quem recebe?Cada assinante conectadoUm consumidor de cada grupo
Receptor offlinePerde a mensagemLê trabalho retido ou pendente depois
Confirmação e repetiçãoNão integradasXACK e reivindicação feita pelo aplicativo
Escala de trabalhadoresDuplica trabalhoDivide entradas
e históricoIndisponíveisDisponíveis durante a retenção
Latência e complexidadeMenoresUm 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.

Um diagrama envia difusões substituíveis a Pub/Sub, tarefas duráveis a Streams e permite um caminho híbrido.
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.

  1. 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.
  2. Baixe o projeto inicial, crie ambiente virtual isolado e instale dependências fixadas.
  3. Crie o e conceda à identidade de desenvolvimento o acesso necessário.
  4. Conecte por com e use a cadeia normal de credenciais do .
  5. Complete funções para publicar evento, publicar em todos os canais configurados e formatar quadros recebidos.
  6. Execute o listener em thread de segundo plano para manter o Flask responsivo.
  7. Assine canais exatos e padrões, rode a página e confirme mensagens ao vivo.
  8. 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

  1. Pub/Sub atende assinantes ativos; retém entradas para consumo posterior ou coordenado.
  2. Use e grupos quando o precisa de confirmação e repetição gerenciada pelo aplicativo.
  3. XADD acrescenta uma entrada ao .
  4. Pub/Sub serve melhor à difusão de status para clientes conectados.
  5. 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.
  • Use , , privilégio mínimo e reconexão testada.

Referências oficiais

  1. geral do
  2. Criar um aplicativo Python com e
  3. Usar na autenticação do
  4. Padrão Publicador-Assinante — Architecture Center
  5. Diretrizes para trabalhos em segundo plano — Architecture Center
  6. Referência de Redis Pub/Sub
  7. Referência de Redis e grupos de consumidores

Resumo do tópico

Mensageria de produção combina o modelo correto com identidade segura, recuperação explícita, idempotência, retenção limitada e observabilidade.