←

Redis Streams con consumer groups para mensajeria durable

Contexto

Para comunicacion entre microservicios tenia dos opciones clasicas: Redis Pub/Sub (rapido pero si el consumer no esta escuchando, pierde el mensaje) o Kafka (durable pero con una complejidad operacional enorme). En ulfblk-redis necesitaba mensajeria durable con acknowledgment, que funcionara con la misma instancia de Redis que ya usaba para cache.

Lo que aprendi

Redis Streams combina lo mejor de ambos mundos: mensajes persistentes como Kafka, pero usando el Redis que ya tienes. Los consumer groups garantizan que cada mensaje se procesa exactamente por un consumer del grupo, y los mensajes no-acknowledeados se pueden reclamar despues de un timeout.

Publicar mensajes al stream

Los mensajes se agregan con XADD. Cada mensaje tiene un ID auto-generado y un payload de key-value pairs.

from redis.asyncio import Redis


class StreamPublisher:
    def __init__(self, redis: Redis):
        self._redis = redis

    async def publish(
        self,
        stream: str,
        data: dict[str, str],
        maxlen: int = 10000,
    ) -> str:
        """Publicar mensaje al stream. Retorna el message ID."""
        message_id = await self._redis.xadd(
            stream,
            data,
            maxlen=maxlen,  # limitar tamano del stream
        )
        return message_id

Crear el consumer group

El consumer group se crea una sola vez. Si ya existe, se ignora el error.

async def ensure_consumer_group(
    redis: Redis,
    stream: str,
    group: str,
) -> None:
    """Crear grupo si no existe. Idempotente."""
    try:
        await redis.xgroup_create(
            stream,
            group,
            id="0",  # leer desde el inicio del stream
            mkstream=True,  # crear stream si no existe
        )
    except Exception as e:
        if "BUSYGROUP" not in str(e):
            raise  # solo ignorar si el grupo ya existe

Consumer que lee y acknowledea mensajes

Cada consumer en un grupo tiene un nombre unico. XREADGROUP entrega mensajes que ningun otro consumer del grupo ha leido. XACK confirma el procesamiento exitoso.

import asyncio


class StreamConsumer:
    def __init__(
        self,
        redis: Redis,
        stream: str,
        group: str,
        consumer_name: str,
    ):
        self._redis = redis
        self._stream = stream
        self._group = group
        self._consumer = consumer_name

    async def consume(
        self,
        handler,
        batch_size: int = 10,
        block_ms: int = 5000,
    ) -> None:
        """Loop principal del consumer."""
        await ensure_consumer_group(self._redis, self._stream, self._group)

        while True:
            # Leer mensajes nuevos asignados a este consumer
            messages = await self._redis.xreadgroup(
                groupname=self._group,
                consumername=self._consumer,
                streams={self._stream: ">"},  # ">" = solo mensajes nuevos
                count=batch_size,
                block=block_ms,
            )

            if not messages:
                continue

            for stream_name, entries in messages:
                for message_id, data in entries:
                    try:
                        await handler(data)
                        # Acknowledge: mensaje procesado exitosamente
                        await self._redis.xack(
                            self._stream, self._group, message_id
                        )
                    except Exception:
                        # No hacer XACK: el mensaje queda en pending
                        pass

Reclamar mensajes fallidos (pending entries)

Si un consumer muere sin hacer XACK, los mensajes quedan en la Pending Entries List (PEL). Otro consumer puede reclamarlos despues de un timeout.

async def claim_stale_messages(
    redis: Redis,
    stream: str,
    group: str,
    consumer_name: str,
    min_idle_ms: int = 60000,  # 1 minuto sin procesar
    count: int = 10,
) -> list[tuple[str, dict]]:
    """Reclamar mensajes que llevan mas de min_idle_ms sin XACK."""

    # Ver mensajes pendientes
    pending = await redis.xpending_range(
        stream, group, min="-", max="+", count=count
    )

    if not pending:
        return []

    # Filtrar por idle time y reclamar
    stale_ids = [
        entry["message_id"]
        for entry in pending
        if entry["time_since_delivered"] >= min_idle_ms
    ]

    if not stale_ids:
        return []

    claimed = await redis.xclaim(
        stream,
        group,
        consumer_name,
        min_idle_time=min_idle_ms,
        message_ids=stale_ids,
    )

    return claimed

Ejemplo de uso completo

# Publicar un evento de orden creada
publisher = StreamPublisher(redis)
await publisher.publish("orders.events", {
    "type": "order_created",
    "order_id": "abc-123",
    "total": "5000",
})

# Consumer que procesa eventos
consumer = StreamConsumer(redis, "orders.events", "billing-group", "billing-1")

async def handle_order_event(data: dict) -> None:
    if data.get("type") == "order_created":
        await create_invoice(data["order_id"], int(data["total"]))

await consumer.consume(handle_order_event)

Por que Streams sobre Pub/Sub o Kafka

Redis Pub/Sub pierde mensajes si no hay nadie escuchando. Kafka requiere brokers, Zookeeper, particiones -- una operacion compleja para equipos chicos. Redis Streams tiene persistencia, consumer groups, acknowledgment, y claim de mensajes fallidos. Todo en el Redis que ya tienes corriendo para cache.

Referencia