Lesson 18 / 25

Partitioning Streams and Ordering

Scale streams beyond one shard while keeping per-key ordering.

One stream, one shard

A stream is a single Redis key, so in a cluster it lives on one shard, and its throughput is limited by that shard's CPU and network. To scale, partition events across several streams, for example events:orders:{0} through events:orders:{15}, choosing the partition by hashing a key such as the customer or order ID. Events with the same key always go to the same stream, which keeps per-key ordering while spreading load across shards. Consumers then read all partitions, usually with one consumer group per partition or a consumer that reads several streams in one XREADGROUP call (only possible when those streams share a hash slot; across slots you read them separately). The trade-off mirrors Kafka partitions: more partitions mean more parallelism but more bookkeeping, and changing the partition count later reshuffles keys, so start with enough. Use hash tags when a stream and related keys (such as a dedupe set) must share a slot for multi-key operations or Lua scripts.

Choosing a partition stream by key

The same order always maps to the same stream, preserving its event order.

import zlib

PARTITIONS = 16

def stream_for(order_id: str) -> str:
    p = zlib.crc32(order_id.encode()) % PARTITIONS
    return f"events:orders:{{{p}}}"           # braces form a hash tag, e.g. events:orders:{7}

def publish(event):
    r.xadd(stream_for(event["orderId"]), event, maxlen=200_000, approximate=True)

# consumers: one worker (or several, via the consumer group) per partition stream
for p in range(PARTITIONS):
    ensure_group(f"events:orders:{{{p}}}", "billing")

Order is per stream, not global

Partitioning keeps each order's events in sequence but loses ordering between different orders. That is usually fine; if you think you need global ordering, check whether you really do.

Quick check: Why partition events across several streams in a Redis Cluster?

  • A single stream lives on one shard, so partitioning spreads load across shards
  • Streams cannot hold more than 100 entries
  • To remove the need for consumer groups
  • To make IDs shorter
Answer

A single stream lives on one shard, so partitioning spreads load across shards — Each stream is one key on one shard; multiple streams use multiple shards.