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())
🛡️ Enterprise Resilience Features in WRedis
-
Pending Entries List (PEL) Recovery: Built-in claiming (
XCLAIM) allows idle or crashed consumer messages to be re-routed to healthy workers automatically. - Memory Governor (MAXLEN ~): Protects Redis RAM from unbounded log growth with approximate stream trimming.
-
Async Native: Fully compatible with
asyncio,FastAPI, andwpipebackground workers.
Top comments (0)