DEV Community

William Rodriguez
William Rodriguez

Posted on

Orquestación reactiva: Redis Streams con grupos de consumidores.

Deja de usar Pub/Sub crudo para procesos críticos donde no puedes perder datos. wpipe-steps aporta pasos nativos de Redis Streams con grupos de consumidores para flujos reactivos resilientes.

Así se utiliza Redis Streams & Consumer Groups en un pipeline de producción con wpipe-steps:

from wpipe import Pipeline
from wpipe_steps.database.redis.streams import redis_stream_add_sync

pipeline = Pipeline(pipeline_name="order_event_dispatcher")

pipeline.set_steps([
    redis_stream_add_sync.as_step(
        name="publish_order_event",
        stream_key="events:orders",
        fields={"order_id": "ORD-9912", "amount": 149.50, "status": "PAID"},
        response_key="stream_message_id"
    )
])

result = pipeline.run({})
print(f"Event published to Redis Stream with ID: {result['stream_message_id']}")
Enter fullscreen mode Exit fullscreen mode

Por qué wpipe-steps marca la diferencia:

  • 196 pasos catalogados que cubren Redis, ClickHouse, MySQL, WAF, S3, Docker e IA local con HuggingFace.
  • Carga perezosa con arranques instantáneos en menos de 100ms.
  • Interfaz factory limpia y consistente con .as_step().

¡Visita el repositorio completo en GitHub!

Autor: William Steve Rodríguez Villamizar (Wisrovi)

Top comments (0)