DEV Community

William Rodriguez
William Rodriguez

Posted on

High-Throughput Event Streaming: Redis Streams & Consumer Groups in WRedis

High-Throughput Event Streaming: Redis Streams & Consumer Groups in WRedis

Day 10 of the Wisrovi Open Source Architecture Series.

When building decoupled distributed architectures, engineering teams often face a dilemma: Apache Kafka brings heavy operational overhead (JVM, ZooKeeper/KRaft, persistent partition state), while classic Redis Pub/Sub offers zero persistence (at-most-once delivery, dropped messages on consumer disconnection).

Redis Streams bridge this gap by providing append-only log structures with consumer group semantics, message acknowledgment (XACK), and fault tolerance directly inside Redis memory.

wredis provides a production-grade, asynchronous wrapper around Redis Streams that removes boilerplate and ensures deterministic stream processing.


🏗️ Architectural Comparison: Streaming Paradigms

Dimension Classic Redis Pub/Sub Apache Kafka wredis Streams
Delivery Guarantee At most once (fire-and-forget) At least once / Exactly once At least once (with XACK)
Consumer Offsets None (instantaneous) Stored in internal topics Consumer Groups (XREADGROUP)
Message Backlog Dropped if consumer is down Persisted to disk Persisted in Redis Stream (Capped)
Operational Cost Minimal High (dedicated cluster) Zero Additional Infra (Uses existing Redis)
Throughput Latency Sub-millisecond Single-digit ms Sub-millisecond

💻 Practical Implementation: Consumer Groups in WRedis

Here is how you initialize consumer groups, stream high-throughput events, and safely acknowledge processed messages:

import asyncio
from wredis import RedisStreamClient

async def run_event_pipeline():
    # Connect to sovereign Redis cluster
    client = RedisStreamClient(host="localhost", port=6379, db=0)
    stream_key = "telemetry:sensor_events"
    group_name = "analytics_processors"
    consumer_name = "worker_node_01"

    # Ensure Consumer Group exists without raising if already present
    await client.ensure_consumer_group(
        stream=stream_key,
        group=group_name,
        create_stream_if_missing=True
    )

    # 1. Producer: Publish event to stream with auto-trimming (MAXLEN)
    event_id = await client.add_event(
        stream=stream_key,
        fields={"device_id": "sensor_42", "temp_c": "28.4", "vibration": "0.012"},
        maxlen=100_000,
        approximate=True
    )
    print(f"Produced event {event_id} to stream {stream_key}")

    # 2. Consumer: Read unacknowledged messages assigned to this consumer
    events = await client.read_group(
        stream=stream_key,
        group=group_name,
        consumer=consumer_name,
        count=10,
        block_ms=2000
    )

    for msg_id, payload in events:
        print(f"Processing message {msg_id}: {payload}")
        # Business logic here...

        # 3. Acknowledge message completion
        await client.ack_event(stream=stream_key, group=group_name, message_id=msg_id)
        print(f"Acknowledged {msg_id}")

    await client.close()

if __name__ == "__main__":
    asyncio.run(run_event_pipeline())
Enter fullscreen mode Exit fullscreen mode

🛡️ Enterprise Resilience Features in WRedis

  1. Pending Entries List (PEL) Recovery: Built-in claiming (XCLAIM) allows idle or crashed consumer messages to be re-routed to healthy workers automatically.
  2. Memory Governor (MAXLEN ~): Protects Redis RAM from unbounded log growth with approximate stream trimming.
  3. Async Native: Fully compatible with asyncio, FastAPI, and wpipe background workers.

redis #python #architecture #eventdriven #backend

Top comments (0)