Using Redis Streams Consumer Groups for Reliable IoT Telemetry Ingestion
Learn how Redis Streams with consumer groups give you durable, exactly‑once processing for high‑throughput event streams, with a Python example and production‑tuning tips.
25 Dec 2025, 13:27 UTC

Problem: losing telemetry is not an option
Imagine a fleet of thousands of sensors pushing temperature, vibration, and GPS readings at tens of thousands of events per second. Your ingestion pipeline must guarantee that every reading is processed at least once, and ideally exactly once, even if a consumer crashes or is restarted. Simple Redis Pub/Sub drops messages when a subscriber is offline, and using plain lists requires manual bookkeeping that is error‑prone.
Thesis: Redis Streams + consumer groups give you built‑in durability and replay
Introduced in Redis 5.0, a Stream is an append‑only log where each entry gets a unique ID. Consumer groups layer a coordinated work‑sharing mechanism on top of that log: each consumer in the group reads new entries, tracks what it has processed, and acknowledges completion with XACK. Unacknowledged entries stay in a pending‑entries list (PEL) and can be reclaimed after a timeout, providing exactly‑once semantics within a single Redis instance.
Worked example: Python producer and two consumers
The following snippet assumes you have a local Redis server running on the default port (6379) and have installed the redis-py library (pip install redis). Run the producer in one terminal and the consumers in separate terminals.
1. Producer – add sensor readings
import redis, time, uuid
r = redis.Redis(host='localhost', port=6379, db=0)
stream = 'iot:telemetry'
# Ensure the stream exists (XADD creates it automatically)
for i in range(10000):
sensor_id = f'sensor-{uuid.uuid4().hex[:8]}'
payload = {
'temp': round(20 + 10 * (i % 3), 2),
'vib': i % 5,
'ts': int(time.time())
}
r.xadd(stream, payload)
if i % 1000 == 0:
print(f'Added {i} entries')
print('Producer finished')
Run: python producer.py. This creates the stream iot:telemetry and appends 10 000 entries.
2. Consumer – join a group and process
import redis, time
r = redis.Redis(host='localhost', port=6379, db=0)
stream = 'iot:telemetry'
group = 'processors'
consumer = f'worker-{int(time.time())}'
# Create the group if it does not exist; start reading from the beginning
try:
r.xgroup_create(stream, group, id='0', mkstream=True)
except redis.ResponseError as e:
if 'BUSYGROUP' not in str(e):
raise
print(f'{consumer} joined group {group}')
while True:
resp = r.xreadgroup(group, consumer, {stream: '>'}, count=10, block=2000)
if not resp:
continue
for _, entries in resp:
for entry_id, fields in entries:
# Simulate processing
print(f'{consumer} processing {entry_id}: {fields}')
# Acknowledge after successful processing
r.xack(stream, group, entry_id)
Start two instances: python consumer.py & twice. Each consumer will receive a subset of the new entries (> means “only deliver unacknowledged messages”). When you stop one consumer with Ctrl+C, the other continues processing the remaining stream.
3. Replay after a crash
To simulate a crash, kill a consumer while it has fetched messages but before it calls XACK. Those messages remain in the group’s pending‑entries list. You can inspect them with:
redis-cli XPENDING iot:telemetry processorsThe output shows the entry IDs, the consumer that last delivered them, and idle time. To have another consumer take over, run:
redis-cli XCLAIM iot:telemetry processors worker-2 0 0 1629345678901-0Replace the ID with the first pending entry from the list. The claimed consumer can now process and acknowledge those messages, guaranteeing none are lost.
Trade‑offs and limitations
- Memory usage. Each stream entry adds overhead (field‑value pairs, ID, and internal structures). Compared with a simple list, a stream can consume several times more memory for the same payload. Use
XTRIMto limit growth:
redis-cli XTRIM iot:telemetry MAXLEN 5000This keeps only the most recent 5 000 entries, discarding older ones. Adjust the limit based on your replay window and available RAM.
- Persistence. Consumer‑group state (the last delivered ID per consumer and the PEL) is stored in Redis memory only. If Redis restarts without persistence, the group state is lost and consumers may reprocess from the beginning or skip messages. Enable AOF (
appendonly yes) or periodic RDB snapshots to survive restarts. - Single‑node guarantee. The exactly‑once promise holds within a standalone Redis instance. In a Redis Cluster, streams are sharded by hash slot; you must ensure all producers and consumers target the same shard, and resharding does not move stream entries. For multi‑node durability, consider using Redis Enterprise’s active‑active geo‑replication or a dedicated stream‑aware broker like Apache Kafka.
Actionable recommendations for production
- Verify version and persistence. Run
redis-cli INFO serverand confirmredis_version≥ "5.0". Ensureappendonly yesis set or schedule regular RDB saves. - Monitor stream size. Periodically check
MEMORY USAGE iot:telemetryorXLEN iot:telemetry. If growth exceeds your budget, schedule anXTRIM job (e.g., via cron) to retain only the needed window. - Tune consumer latency. The
blockargument inXREADGROUPlets you balance latency vs. CPU usage. Start with 2000 ms (2 seconds) and adjust based on your SLA. - Handle processing failures. If a consumer crashes after processing but before
XACK, the message stays pending. Implement a dead‑letter mechanism: after a configurable idle time (checked viaXPENDING), move the entry to a separate stream for manual inspection. - Test replay. In a staging environment, kill a consumer mid‑batch, verify pending entries, then start a new consumer and confirm it processes those entries exactly once.
Closing
Redis Streams with consumer groups give you a lightweight, durable queue that fits naturally into existing Redis‑based stacks. By combining the append‑only log model with automatic acknowledgment and pending‑entry tracking, you gain exactly‑once processing without adding a separate messaging system. Keep an eye on memory, enable persistence, and trim old data to keep the solution cost‑effective for high‑volume IoT telemetry or any event‑driven workload.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.