DEV Community

William Rodriguez
William Rodriguez

Posted on

Event-driven orchestration: Redis Streams with consumer groups.

Stop using raw Redis Pub/Sub for workflows that cannot afford message loss. wpipe-steps provides native Redis Streams steps with consumer group support for resilient event-driven architectures.

Here is how you use Redis Streams & Consumer Groups in a production pipeline with 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

Why developers love wpipe-steps:

  • 196 cataloged steps covering Redis, ClickHouse, MySQL, WAF, S3, Docker, and local HuggingFace AI.
  • Lazy-loading imports for instant sub-100ms startup times.
  • Clean .as_step() factory interface.

Explore the complete modular catalog on GitHub!

Author: William Steve Rodríguez Villamizar (Wisrovi)

Top comments (1)

Collapse
 
william_rodriguez_65a5898 profile image
William Rodriguez •

En sistemas distribuidos reactivos con Kafka, la frontera entre control de flujo y resiliencia ante caídas es sutil. El uso de patrones declarativos con commits atómicos permite que la lógica de reintento no degrade el rebalanceo de particiones.

¿Prefieren delegar el manejo de backpressure y offset commits a middleware declarativo o gestionarlo manualmente por handler?