Azure Managed Redis: Pub/Sub, Streams, eventos y colas de tareas
Volver a la ruta AI-200
AI-200Capítulo 16

Preparación para la Certificación Microsoft AI-200

Azure Managed Redis: Pub/Sub, Streams, eventos y colas de tareas

Difunde eventos de IA en vivo, coordina trabajo duradero con Redis Streams, recupera entregas abandonadas, escala grupos de consumidores y combina ambos modelos de forma segura.

Tiempo de estudio sugerido: 110 minutos • Nivel intermedio • Reescritura original completa con versión resumida de cada tema, evaluación comentada y laboratorio guiado en Python

Escudo neón Microsoft Certified AI-200 con canales Pub/Sub, Streams, grupos de consumidores, colas de tareas y recuperación en Azure Managed Redis

1. Desacoplar una canalización de IA en tiempo real

Imagina una plataforma jurídica que ejecuta OCR, reconocimiento de entidades, clasificación y generación de embeddings para cientos de documentos simultáneos. Una sincrónica mantendría cada solicitud abierta durante segundos; el sondeo y los reintentos personalizados serían frágiles. La mensajería permite responder rápido a la carga, escalar cada etapa por separado y enviar progreso a paneles en vivo.

y ofrecen Pub/Sub y Redis . Pub/Sub distribuye una notificación transitoria a los oyentes conectados. conserva el trabajo para que consumidores coordinados lo procesen y confirmen. Los ejemplos usan Python y redis-py; consulta las firmas de la versión instalada porque su evoluciona.

  • Difundir un evento a varios servicios de IA mediante canales y patrones.
  • Crear una cola duradera con y grupos de consumidores.
  • Recuperar de forma explícita el trabajo pendiente tras un error.
  • Elegir difusión, distribución coordinada o arquitectura híbrida.
  • Realizar un ejercicio Flask con y la CLI de .
Una API de IA envía eventos transitorios por Redis Pub/Sub y trabajos duraderos por Redis Streams a observadores independientes y trabajadores coordinados.
Una entrada puede generar a la vez una notificación en vivo y una unidad de trabajo duradera.

Resumen del tema

Separa la solicitud rápida del procesamiento asíncrono y usa Pub/Sub para difusión en vivo y para trabajo retenido.

2. Distinguir eventos, mensajes y tareas

La intención determina el modelo.
CargaSignificadoExpectativa
EventoHecho que ya ocurrió, como model_updatedCero, una o varias reacciones independientes
NotificaciónActualización breve, como prediction_readySolo los oyentes conectados pueden necesitarla
Tarea o comandoSolicitud de trabajo, como analyze_documentUn trabajador debe completarla, reintentar o registrar el fallo

El publicador no debería conocer cada punto de conexión receptor. El intermediario desacopla implementación y escala; un identificador de correlación, tipo, versión de esquema, marca de tiempo y referencia a cargas grandes conservan trazabilidad. Mantén documentos y artefactos fuera del mensaje y transporta identificadores.

Resumen del tema

Los eventos anuncian hechos y las tareas solicitan trabajo; modela la carga y su entrega antes de elegir la función Redis.

3. Comprender la difusión por canales Redis Pub/Sub

El publicador ejecuta PUBLISH con canal y carga. Redis envía inmediatamente una copia a cada suscriptor activo sin que el emisor conozca a los oyentes. Así, los servicios de sentimiento, intención y contexto reaccionan en paralelo a un evento de conversación nueva.

El desacoplamiento ofrece baja latencia y gran rendimiento. Es apropiado para coordinación en tiempo real, invalidación de caché, avisos de actualización de modelos o embeddings, estado de entrenamiento, predicciones y telemetría útil mientras el consumidor está conectado.

Un publicador envía a un canal Redis Pub/Sub que copia el evento a cuatro suscriptores conectados; un suscriptor sin conexión lo pierde.
Pub/Sub es un bus de difusión para oyentes conectados, no una cola duradera.

Resumen del tema

PUBLISH crea una difusión uno a muchos: cada suscriptor conectado recibe el evento de manera independiente.

4. Diseñar canales y suscripciones por patrón

Los canales son cadenas. Un predecible muestra dominio, asunto y evento, como ai:models:updated. Añade inquilino o entidad solo cuando el aislamiento o la entrega dirigida lo exijan; los canales ilimitados por usuario son difíciles de operar.

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

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

SUBSCRIBE nombra canales exactos. PSUBSCRIBE acepta patrones glob y entrega un pmessage con patrón, canal real y carga. Los patrones amplios facilitan descubrimiento, pero aumentan tráfico y pueden exponer eventos no relacionados; autoriza y documenta los espacios de nombres.

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"],
        )

Resumen del tema

Usa espacios estables, suscripciones exactas para tráfico estrecho y patrones limitados para familias relacionadas.

5. Contemplar entrega como máximo una vez y contrapresión

  • Sin persistencia: el mensaje no queda disponible para reproducción.
  • Entrega como máximo una vez: un suscriptor desconectado, reiniciándose o con error pierde el evento.
  • Sin confirmación: el recuento devuelto no demuestra que el proceso de negocio terminó.
  • Sin contrapresión de cola: el suscriptor lento debe almacenar, descartar o desconectar en su propio proceso.
  • Sin historial: orden, auditoría y requieren otro almacén o modelo.

Estas propiedades ayudan cuando el evento es reemplazable o efímero. No sirven si perder un documento o cobrar dos veces importa. No conviertas Pub/Sub en cola asignando la misma suscripción a todos los trabajadores: todos recibirán la tarea y multiplicarán el coste.

Resumen del tema

Pub/Sub prioriza velocidad y fan-out sobre durabilidad, confirmación, y reparto coordinado de carga.

6. Aplicar Pub/Sub a los escenarios correctos de IA

Escenarios de difusión.
EscenarioPor qué sirve fan-out
Invalidación de modelo o embeddingCada instancia limpia sus datos locales obsoletos
Configuración y marcas de característicasTodos los servicios cargan el mismo cambio
Predicción preparada, paneles y telemetría reaccionan por separado
Métricas de IAAlertas, paneles y registros consumen la misma señal de modo distinto
Acciones heterogéneasAnálisis, facturación y recomendaciones actúan de manera diferente tras una interacción

El error de un suscriptor no debe bloquear al publicador ni al resto. Si el hecho también requiere una acción garantizada, persiste la tarea aparte —por ejemplo, en un — y usa Pub/Sub solo para la vista en vivo.

Resumen del tema

Elige Pub/Sub cuando todos los servicios conectados deben ver el mismo evento reemplazable y reaccionar de forma independiente.

7. Publicar y escuchar con redis-py

La conexión presupone el puerto 10000 de y un proveedor configurado. En producción, ejecuta el listener bloqueante en un trabajador o hilo dedicado, controla reconexiones, valida esquemas y cierra bien la suscripción.

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"])

Comprueba event["type"] porque también llegan marcos de confirmación. decode_responses=True convierte canales y cargas en cadenas; desactívalo para binarios. El entero de publish() indica cuántos clientes recibieron el marco, no cuántos completaron la acción.

Resumen del tema

Usa publish(), pubsub(), subscribe(), psubscribe() y listen() con tratamiento explícito, conexión segura y ciclo resiliente.

8. Reenviar eventos Redis al navegador

El navegador no suele conectarse directamente a Redis. Un servicio FastAPI o Flask de confianza se suscribe en segundo plano y reenvía eventos autorizados por . Debe asignar usuarios a temas permitidos, limitar colas por , cerrar conexiones inactivas y evitar fugas entre inquilinos.

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"]})

Una suscripción Redis por navegador puede agotar conexiones. Un de producción suele compartir una por proceso, desmultiplicar hacia locales y publicar solo estado compacto; el resultado duradero queda detrás de autenticadas.

Resumen del tema

El convierte fan-out Redis en actualizaciones del navegador con autorización, búfer y límites.

9. Modelar trabajo duradero como Redis

Un Redis es una secuencia append-only de campos y valores. XADD agrega una tarea y devuelve un identificador ordenado en el tiempo, como 1699980000000-0. A diferencia de Pub/Sub, la entrada permanece hasta eliminarla o recortarla. La encola inferencia y devuelve el ID mientras los trabajadores procesan.

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"}

sirve para inferencias, canalizaciones extraer-analizar-resumir-incorporar, operaciones largas, historial y recuperación. El ID es una buena clave de correlación e idempotencia, pero no convierte el efecto de negocio en exactamente una vez.

Resumen del tema

XADD crea un registro ordenado y retenido, por lo que la puede responder antes del procesamiento de IA.

10. Distribuir tareas con grupos de consumidores

XGROUP CREATE establece el grupo y MKSTREAM puede crear el . Los trabajadores comparten el nombre del grupo, pero usan nombres de consumidor únicos y sensibles a mayúsculas. XREADGROUP con > solicita entradas no entregadas todavía en ese grupo. Redis reparte entradas nuevas entre lectores activos y permite escalar sin balanceo en el 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 quita la entrega completada de la Pending Entries List (PEL) del grupo, pero no necesariamente elimina la entrada. Leer con un ID como 0 accede a pendientes del consumidor, no a trabajo nuevo. Varios grupos consumen el mismo de forma independiente.

Resumen del tema

Los grupos reparten entradas nuevas entre trabajadores únicos y XACK registra el procesamiento correcto.

11. Recuperar de forma explícita las entregas con error

Si un trabajador muere tras recibir una entrada y antes de XACK, Redis la conserva en la PEL. No la reasigna ni reintenta automáticamente. La aplicación debe consultar XPENDING y reclamar trabajo inactivo con XCLAIM o XAUTOCLAIM. Un consumidor reiniciado también puede leer su historial pendiente antes de pedir novedades.

# 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)

La recuperación produce comportamiento al menos una vez: otro trabajador puede reclamar un elemento mientras el primero termina. Haz idempotentes los efectos mediante el ID del o una clave de negocio, limita reintentos y mueve tareas problemáticas a un de mensajes fallidos con diagnóstico.

Resumen del tema

La PEL conserva trabajo no confirmado, pero la aplicación debe reclamar entradas inactivas y tolerar nuevas entregas.

12. Supervisar y limitar la retención del

XINFO muestra longitud e IDs; XINFO GROUPS presenta retraso y pendientes; XINFO CONSUMERS revela consumidores, inactividad y propiedad. Alerta ante retraso creciente, pendientes antiguas, demasiados intentos, rotación de trabajadores, memoria, latencia y errores.

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,
)
  • Recorta por longitud máxima con XADD MAXLEN o XTRIM; el modo aproximado cuesta menos.
  • Mantén cargas pequeñas y documentos o embeddings en almacenamiento duradero.
  • Define retención suficiente para auditoría y recuperación; un ilimitado no caduca solo.
  • Persistencia y confirmación añaden latencia y código frente a Pub/Sub.
  • Prueba apagado, reconexión, mensajes problemáticos y recuperación bajo carga.

Resumen del tema

Observa , grupos y consumidores y limita entradas para que la fiabilidad no provoque crecimiento sin control.

13. Elegir difusión o distribución coordinada

Guía de decisión.
RequisitoPub/Sub con grupo
¿Quién recibe?Cada suscriptor conectadoUn consumidor de cada grupo
Receptor sin conexiónPierde el mensajeLee trabajo retenido o pendiente después
Confirmación y reintentoNo integradosXACK y reclamación por la aplicación
Escala de trabajadoresDuplica trabajoReparte entradas
e historialNo disponiblesDisponibles durante la retención
Latencia y complejidadMenoresAlgo mayores

Usa Pub/Sub para invalidación, configuración, métricas y avisos descartables. Usa para trabajos, inferencia y canalizaciones recuperables. Si dominan transacciones, enrutamiento avanzado, mensajes fallidos o garantías entre sistemas, compara un servicio de mensajería de especializado.

Un diagrama dirige difusiones reemplazables a Pub/Sub, tareas duraderas a Streams y permite un camino híbrido.
La intención de entrega, no solo el rendimiento, determina el mecanismo.

Resumen del tema

Difunde hechos con Pub/Sub, coordina trabajo duradero con y evalúa un broker dedicado para requisitos más amplios.

14. Combinar y Pub/Sub en una arquitectura

Muchos sistemas necesitan ambos. Al recibir un documento, escribe el trabajo duradero en el y publica un evento efímero. Un miembro del grupo procesa cada entrega, mientras , análisis y supervisión ven las difusiones. Etapas posteriores pueden escribir otros y publicar progreso.

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,
}))

Las dos escrituras no forman una transacción de negocio por defecto. Si perder el aviso es aceptable, trátalo como mejor esfuerzo. Si deben ser coherentes, usa un outbox u otro registro duradero. Repite el mismo ID de correlación en , Pub/Sub, registros y resultados.

Resumen del tema

El diseño híbrido usa como trabajo oficial y Pub/Sub para observabilidad efímera, unidos por correlación y reglas de coherencia.

15. Laboratorio guiado: publicar y suscribirse en

El ejercicio de origen estima 30–40 minutos y crea una página Flask con Python para publicar y suscribirse en tiempo real. Usa un recurso desechable, no confirmes credenciales en el código y elimina el entorno al terminar.

  1. Prepara una suscripción de , , Python 3.12 o posterior, la CLI de más reciente y la extensión redisenterprise instalada con az extension add --name redisenterprise.
  2. Descarga el proyecto inicial, crea un entorno virtual aislado e instala dependencias fijadas.
  3. Crea y concede a la identidad de desarrollo el acceso necesario.
  4. Conecta por con mediante la cadena normal de credenciales de .
  5. Completa funciones para publicar un evento, publicar en todos los canales configurados y dar formato a los marcos.
  6. Ejecuta el listener en un hilo de fondo para mantener Flask disponible.
  7. Suscríbete a canales exactos y patrones, ejecuta la página y verifica mensajes en vivo.
  8. Prueba una desconexión para observar que Pub/Sub no reproduce lo perdido y elimina el recurso.

Resumen del tema

El laboratorio valida aprovisionamiento seguro, escucha de fondo, suscripciones exactas y por patrón, semántica y limpieza.

16. Revisión de la evaluación y lista de producción

  1. Pub/Sub atiende a suscriptores activos; retiene entradas para consumo posterior o coordinado.
  2. Usa y grupos cuando la canalización requiere confirmación y reintento administrado por la aplicación.
  3. XADD agrega una entrada al .
  4. Pub/Sub encaja con la difusión de estado a clientes conectados.
  5. XREADGROUP con > reparte entradas nuevas; la idempotencia protege ante entregas repetidas.
  • Versiona esquemas e incluye identificadores de correlación e idempotencia.
  • Autoriza espacios de canales y deja cargas grandes fuera del mensaje.
  • Da un nombre único a cada consumidor.
  • Confirma solo después del efecto de negocio.
  • Supervisa antigüedad de PEL, retraso, reintentos, memoria, latencia y desconexiones.
  • Implementa XAUTOCLAIM o XPENDING con XCLAIM; Redis no recupera solo.
  • Recorta y trata tareas problemáticas con una estrategia de mensajes fallidos.
  • Usa , , privilegio mínimo y reconexión probada.

Referencias oficiales

  1. Introducción a
  2. Crear una aplicación Python con y
  3. Usar para autenticar
  4. Patrón Publicador-Suscriptor — Architecture Center
  5. Guía para trabajos en segundo plano — Architecture Center
  6. Referencia de Redis Pub/Sub
  7. Referencia de Redis y grupos de consumidores

Resumen del tema

La mensajería de producción combina el modelo correcto con identidad segura, recuperación explícita, idempotencia, retención limitada y observabilidad.