Range sharding splits a keyspace into contiguous, sorted intervals and assigns each interval to a node. Keys from a to f live on one shard, f to m on another, and so on. It is the partitioning scheme behind Bigtable tablets, HBase regions, CockroachDB ranges, TiKV regions and MongoDB ranged sharding, and it exists for one reason: it keeps neighbouring keys together. A query for every order of customer 42 in September, or every event between two timestamps, touches one shard or a few adjacent ones, where hash sharding would have to ask every shard.

That locality comes with a price. Ranges grow unevenly, so the system must split them; traffic concentrates on whichever range holds the newest keys, so it must detect and spread hotspots; and the map from range to node changes constantly, so every client and router must cope with being wrong. This article covers the range map and its lookup code, how scans cross boundaries, how to pick split points by size and by load, merges and rebalancing, the monotonic-key trap and how to design keys around it, stale routing, and a worked example of an event store, ending with a checklist.

Keyspace, sorted[ "" ............................................................. +inf )R1 ["", c)node 1, gen 4R2 [c, k)node 2, gen 7R3 [k, p)node 1, gen 2R4 [p, +inf)node 3, gen 9, hotR4a [p, t)node 3, gen 10R4b [t, +inf)moved to node 2Meta range: sorted list of descriptorsstart key, end key, replicas, generationRouter cachebisect on start keys, retry on stale gensplit at twatchroute
A range-sharded keyspace: descriptors live in a meta range, routers cache them and bisect on start keys, and a hot range is split by load and half of it moved.

The range map and lookup

Every range-sharded system has the same core data structure: a sorted list of range descriptors, each with an inclusive start key, an exclusive end key, the replicas that hold it and a generation number that increases whenever the range changes. The first range starts at the empty key and the last ends at infinity, so every possible key belongs to exactly one range. The list itself is stored somewhere durable, typically in a special meta range of the same database, replicated with consensus. Routers and clients cache it and look up a key by binary search on the start keys.

import bisect
from dataclasses import dataclass

@dataclass
class Range:
    start: bytes            # inclusive
    end: bytes | None       # exclusive; None means +infinity
    node: str
    gen: int = 1            # bumped on every split, merge or move

class RangeMap:
    def __init__(self, ranges):
        self.ranges = sorted(ranges, key=lambda r: r.start)
        assert self.ranges[0].start == b"", "first range must start at the empty key"
        self.starts = [r.start for r in self.ranges]

    def lookup(self, key: bytes) -> Range:
        i = bisect.bisect_right(self.starts, key) - 1
        return self.ranges[i]

    def scan(self, lo: bytes, hi: bytes):
        """Yield (range, sub_lo, sub_hi) pieces that exactly cover [lo, hi)."""
        i = bisect.bisect_right(self.starts, lo) - 1
        while i < len(self.ranges) and self.ranges[i].start < hi:
            r = self.ranges[i]
            sub_hi = hi if r.end is None else min(hi, r.end)
            yield r, max(lo, r.start), sub_hi
            i += 1

    def split(self, at: bytes):
        i = bisect.bisect_right(self.starts, at) - 1
        r = self.ranges[i]
        if at == r.start:
            raise ValueError("split key must be strictly inside the range")
        left = Range(r.start, at, r.node, r.gen + 1)
        right = Range(at, r.end, r.node, r.gen + 1)
        self.ranges[i:i + 1] = [left, right]
        self.starts[i:i + 1] = [left.start, right.start]
        return left, right

Two things in this code are worth copying. Lookup is bisect_right minus one, which finds the last range whose start is at or before the key; with start keys inclusive this is exactly right, including a key equal to a boundary. And the generation number is part of the descriptor, because the router will later need to tell whether the descriptor it used is still current. Both functions were exercised against a brute-force linear scan over random keys and random split sequences before being published here.

The range map and lookup

Every range-sharded system has the same core data structure: a sorted list of range descriptors, each with an inclusive start key, an exclusive end key, the replicas that hold it and a generation number that increases whenever the range changes. The first range starts at the empty key and the last ends at infinity, so every possible key belongs to exactly one range. The list itself is stored somewhere durable, typically in a special meta range of the same database, replicated with consensus. Routers and clients cache it and look up a key by binary search on the start keys.

import bisect
from dataclasses import dataclass

@dataclass
class Range:
    start: bytes            # inclusive
    end: bytes | None       # exclusive; None means +infinity
    node: str
    gen: int = 1            # bumped on every split, merge or move

class RangeMap:
    def __init__(self, ranges):
        self.ranges = sorted(ranges, key=lambda r: r.start)
        assert self.ranges[0].start == b"", "first range must start at the empty key"
        self.starts = [r.start for r in self.ranges]

    def lookup(self, key: bytes) -> Range:
        i = bisect.bisect_right(self.starts, key) - 1
        return self.ranges[i]

    def scan(self, lo: bytes, hi: bytes):
        """Yield (range, sub_lo, sub_hi) pieces that exactly cover [lo, hi)."""
        i = bisect.bisect_right(self.starts, lo) - 1
        while i < len(self.ranges) and self.ranges[i].start < hi:
            r = self.ranges[i]
            sub_hi = hi if r.end is None else min(hi, r.end)
            yield r, max(lo, r.start), sub_hi
            i += 1

    def split(self, at: bytes):
        i = bisect.bisect_right(self.starts, at) - 1
        r = self.ranges[i]
        if at == r.start:
            raise ValueError("split key must be strictly inside the range")
        left = Range(r.start, at, r.node, r.gen + 1)
        right = Range(at, r.end, r.node, r.gen + 1)
        self.ranges[i:i + 1] = [left, right]
        self.starts[i:i + 1] = [left.start, right.start]
        return left, right

Two things in this code are worth copying. Lookup is bisect_right minus one, which finds the last range whose start is at or before the key; with start keys inclusive this is exactly right, including a key equal to a boundary. And the generation number is part of the descriptor, because the router will later need to tell whether the descriptor it used is still current. Both functions were exercised against a brute-force linear scan over random keys and random split sequences before being published here.

Scans across range boundaries

Point lookups are where range sharding is merely competitive; range scans are where it wins. A scan over [lo, hi) is decomposed by the scan generator into one sub-scan per overlapping range, each clipped to that range's bounds. The sub-scans can run in parallel for throughput or sequentially for a paginated, ordered result; sequential is what lets ORDER BY key LIMIT 100 stop after one shard.

The subtle part is that the map can change while a scan is running. If range R splits into R-left and R-right after you computed the plan, your sub-scan request to R arrives at a node that now owns only part of it. Correct systems handle this in the storage node, not the client: the node checks the request's descriptor generation, and if it is stale returns an error carrying the new descriptors. The client then re-plans only the unfinished part, resuming from the last key it received. Write your pagination cursor as the last key returned, never as a shard id plus offset, and this resumption is trivial.

Choosing split points: size versus load

Ranges are split for two different reasons, and the split point is chosen differently for each.

Size-based splitting keeps ranges small enough to move, replicate and recover quickly. When a range exceeds a configured size, the node picks a key near the byte midpoint, which storage engines can find cheaply from index blocks, and splits there. This keeps recovery and rebalancing granular but does nothing for throughput: a small range can still be the hottest thing in the cluster.

Load-based splitting targets throughput. The node samples requests to the range, by queries per second or CPU time, and looks for a key that divides the load, not the bytes, roughly in half. The sketch below is the essential logic.

def pick_load_split(samples, min_share=0.2):
    """samples: (key, weight) pairs sampled from requests to one range.
    Returns a split key with at least min_share of the load on each side, or None."""
    if len(samples) < 100:
        return None                       # not enough evidence yet
    samples = sorted(samples)
    total = sum(w for _, w in samples)
    acc = 0.0
    for key, w in samples:
        acc += w
        if acc >= total / 2:
            break                         # key is the weighted median
    left = sum(w for k, w in samples if k < key)    # load that would move left
    if left < min_share * total or total - left < min_share * total:
        return None                       # one key dominates: splitting cannot help
    return key

The rejection rule is the important line. If one key carries most of the load, for example a single celebrity account or a global counter row, no split can divide it, and splitting anyway just creates a stream of tiny ranges that are all equally useless. Real systems also require the imbalance to persist for some time before acting, so that a brief burst does not leave a trail of permanent splits. Note too that splitting alone does not add capacity: both halves start on the same node, and the gain comes only when the rebalancer moves one half elsewhere.

Merges and rebalancing

Splits accumulate. Data is deleted, a burst of load passes, a table is truncated, and you are left with thousands of near-empty ranges, each with its own replication group, heartbeats and metadata entry. Merging adjacent ranges that are both small and cold reverses this. A merge is more delicate than a split because the two ranges may live on different replica sets; the system must first co-locate them, then atomically replace two descriptors with one, bumping the generation so that cached copies of either old descriptor are rejected.

Rebalancing moves whole ranges between nodes to even out disk usage, load and leadership. Because ranges are sized for this, a move is a replica addition, a snapshot transfer and a replica removal. The rebalancer must be rate-limited, since moving a range costs disk and network on both sides, and it must not fight the splitter: a range that was just split for load should have one half moved, not be merged back because both halves now look small.

The monotonic key hotspot

The best known weakness of range sharding is the monotonic key. If the primary key is a timestamp, an auto-increment integer or a time-ordered identifier, every new row sorts after every existing row, so every insert goes to the last range. Splitting the last range moves the old data to a new range and leaves the hot tail where it was; the next insert lands on the tail again. Load-based splitting cannot fix this, and the cluster's write throughput collapses to that of one replica group.

The fixes all change the key, and each one trades away some of the ordering you chose range sharding for:

Key designWrites spread?What you lose
Random id (UUIDv4) as leading columnyes, uniformlytime-ordered scans; also worse index locality
Hash-prefix bucket: (hash(id) % 16, ts, id)yes, across 16 bucketsa time scan becomes 16 parallel scans merged
Natural high-cardinality prefix: (tenant_id, ts)yes, if tenants are many and similarcross-tenant time scans; a giant tenant is still a hotspot
Reversed or bit-reversed sequenceyesany range query on that column
Keep monotonic key, accept one hot rangenonothing, if write rate fits one replica group

The bucketed form is the usual compromise for time series. Writes spread over a fixed number of buckets, and a query for the last hour fans out to exactly that many scans, each still ordered by time, which the client merges. Choose the bucket count to match the write parallelism you need, not the node count, because nodes will change and the bucket count cannot without rewriting keys.

Stale routing and generations

Every router holds a cached copy of the map, and that copy is always slightly wrong. The rule that keeps the system correct is that storage nodes are the authority, caches are hints. Each request carries the range id and generation the router believes in. The node serving it checks that it still holds that range at that generation and that the key lies within its current bounds; if not, it rejects the request with the current descriptor or a hint of where the key went. The router updates its cache and retries, with a retry limit so a corrupt cache cannot cause a loop.

Getting this wrong produces the most dangerous range-sharding bug: a write accepted by a node that no longer owns the key because a split or move completed a moment earlier. The node-side bounds check closes that window. Leases close the related one for reads, where a former leaseholder could serve a stale read; a node may only serve reads while it holds an unexpired lease for the range.

Worked example: an order events store

An order-events store takes 40,000 writes per second at peak and serves two queries: all events for one order, and all events in a time window for an analytics export. A first design keys rows by (event_ts, order_id, seq). In load testing one node sits at full CPU while the rest idle; the hot range is always the last one, and load-based splitting logs that it declined to split because the median key moved past the split point before the move completed. This is the monotonic-key trap.

The second design keys by (order_id, event_ts, seq) with order ids drawn randomly. Writes spread across all ranges, and the per-order query is a single-range scan. The analytics export is now a full scan with a time filter. Because exports run hourly and the team also streams changes to a warehouse, they accept that, and move time-window analytics to the warehouse. If the window query had needed to stay online, the alternative was (bucket, event_ts, order_id, seq) with 32 buckets: 32 merged scans per window query, and per-order lookups through a secondary index. The point is that the key is chosen from the query list, with the write pattern as a hard constraint.

Failure modes

  • Monotonic leading key. One hot range, splits that do not help. Detect it by watching for the hottest range always being the last.
  • Single hot key. No split point exists. Cache, shard the counter into sub-keys, or batch writes.
  • Split storms. Bursty traffic triggers splits that never merge. Require sustained imbalance and enable merging of cold neighbours.
  • Client-side ownership trust. A write lands on a node that just gave up the range. Always enforce bounds and generation on the node.
  • Offset-based pagination. Breaks when ranges split mid-scan. Page by last key.
  • Meta range overload. Cold caches across a fleet restart stampede the meta range. Stagger restarts and prefetch.

What to do next

Keep going with database sharding architecture for the shard map and router tier in general, directory-based sharding for per-entry placement, sharded architecture for sharding a whole service stack, and geo-partitioning for pinning key ranges to regions.

  1. List your queries and check that the ones you care about are prefix scans on the proposed key. If none are, hash sharding is simpler.
  2. Check the write pattern: if the leading key column is monotonic, change it or bucket it before launch.
  3. Implement routing with binary search on inclusive start keys and carry the descriptor generation on every request.
  4. Make storage nodes reject requests whose key or generation does not match, and make clients retry with a bounded count.
  5. Page every scan by last key, never by shard and offset.
  6. Enable both size- and load-based splitting, require sustained imbalance, and turn on merging of small cold neighbours.
  7. Dashboard per-range QPS and size; alert when the same range is hottest for hours, which means a key-design problem rather than a capacity one.
Key takeaway: Range sharding keeps neighbouring keys together so prefix scans touch few shards. Route by binary search over inclusive start keys, carry a generation on every request and let storage nodes reject stale ones. Split by size for manageability and by load for throughput, refuse to split a single hot key, and merge cold neighbours. Above all, never lead a key with a monotonic value unless one range can absorb all writes.