Stackbook Logo
scalingestablished · high operational burden

Partitioning (Sharding)

Also known as: sharding, horizontal-partitioning, data-partitioning

Intent

Distribute data and load across multiple nodes by partitioning on a key, enabling horizontal scale beyond single-node limits.

Problem

Single database/node hits limits: storage, CPU, memory, connections. Vertical scaling has hard ceiling.

Forces

  • Data growth exceeds single-node capacity
  • Query latency must stay low as data grows
  • Need to add capacity without downtime
  • Cross-partition queries are expensive/complex

Solution

✓ When to Use

  • Data size > single node capacity (TB+)
  • Throughput > single node (100k+ QPS)
  • Multi-tenant isolation required
  • Geo-distribution (data residency)

✗ When Not to Use

  • Data fits on one node with room to grow
  • Complex cross-partition queries (joins, transactions)
  • Team not ready for operational complexity

Pros

  • +Horizontal scale: add nodes, not bigger nodes
  • +Fault isolation: one partition down ≠ total outage
  • +Tenant isolation: noisy neighbor contained

Cons

  • Cross-partition queries: scatter-gather, high latency
  • Resharding is complex, risky, long-running
  • Hotspots: skew in key distribution
  • Operational burden: more nodes, more failure domains

Cost Profile

Infrastructure

High — more nodes, coordination overhead

Operational

High — resharding, monitoring per partition

Cognitive

High — partition key choice, query routing

Failure Modes

  • Hot partition: one shard gets 80% load

  • Resharding stall: migration falls behind

  • Cross-shard transaction failure: partial commit

  • Partition count too high: metadata overhead, slow planning

Real-World Examples

Alternatives

  • read-replica
  • vertical-scaling
  • caching
  • archive-old-data

Related Patterns

  • consistent-hashing
  • read-replica
  • tenant-isolation
  • multi-tenancy

Competency Domains

data statescalingeconomics evolutionreliability ops