DEV Community

William Rodriguez
William Rodriguez

Posted on

Ingestión confiable de eventos en Python: Redis Streams y Consumer Groups con WRedis

Cuando tu aplicación supera las capacidades de Redis Pub/Sub tradicional, requieres persistencia en disco, control de offsets por consumidor y distribución horizontal de carga sin la complejidad de desplegar un broker pesado.

WRedis v1.0.0 LTS implementa abstracciones declarativas para Redis Streams (XADD, XREADGROUP) con decoradores automáticos para grupos de consumidores y parada limpia ante señales del sistema.

Este es el Día 03 de la serie técnica WRedis Open Source (MIT, Python 3.10+, 95%+ test coverage).


¿Por qué elegir Redis Streams frente a Pub/Sub?

  1. Historial Persistente: Los eventos se almacenan en el log de Redis y no se pierden si un worker se reinicia.
  2. Consumer Groups: Balanceo automático del procesamiento entre múltiples instancias sin duplicar mensajes.
  3. Pending Entries List (PEL): Recuperación garantizada de mensajes no confirmados tras caídas imprevistas.

Implementación en Producción: Streams con @on_message

1. Ingesta de Eventos Estructurados (add_to_stream)

from wredis.streams import RedisStreamManager

stream_mgr = RedisStreamManager(host="localhost", port=6379)

# Enviar evento estructurado con serialización JSON automática
msg_id = stream_mgr.add_to_stream(
    key="eventos:transacciones",
    data={"tx_id": "TX-9941", "monto": 250.0, "moneda": "EUR"},
    ttl=86400  # Retención del stream por 24h
)
print(f"Evento registrado con ID: {msg_id}")
Enter fullscreen mode Exit fullscreen mode

2. Worker con Consumer Group Declarativo

from wredis.streams import RedisStreamManager

stream_mgr = RedisStreamManager(host="localhost", verbose=False)

# Registro del worker con grupo de consumidores autogestionado
@stream_mgr.on_message(
    stream_name="eventos:transacciones",
    group_name="workers_liquidacion",
    consumer_name="nodo_primario"
)
def liquidar_transaccion(data):
    print(f"Liquidando TX: {data['tx_id']} por {data['monto']} {data['moneda']}")

# Iniciar bucle de consumo multihilo con parada limpia
stream_mgr.wait()
Enter fullscreen mode Exit fullscreen mode

Ventajas de WRedis en Producción

  • Grupos Sin Fricción: Crea los Consumer Groups automáticamente al iniciar sin arrojar excepciones si ya existen.
  • Hilos Gestionados: Los listeners corren en segundo plano con reconexión automática sin bloquear tu aplicación.
  • 95%+ Cobertura de Tests: Validado con 800+ pruebas unitarias, 38 de integración y 19 de estrés concurrente.

Instalación y Recursos

pip install wredis
Enter fullscreen mode Exit fullscreen mode

Autor: William Steve Rodríguez Villamizar (Wisrovi)

Top comments (0)