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.