Partitioning skill

Partitioning (also called sharding) divides a large dataset into smaller subsets called partitions, each stored on a different node.

by wondelai·MIT license·★ 2,235 Stars on the repo·GitHub ↗

Use now

Files of Partitioning

wondelai/main1 file
partitioning.md
Show the full text295 lines

Partitioning

Partitioning (also called sharding) divides a large dataset into smaller subsets called partitions, each stored on a different node. The goal is to spread data and query load evenly across machines, enabling horizontal scaling beyond the capacity of a single node.

Why Partition?

A single database node has hard limits:

  • Storage capacity: A single disk or SSD has a maximum size
  • Write throughput: A single CPU can process a limited number of writes per second
  • Read throughput: Even with caching, a single node can serve a limited number of concurrent reads

Partitioning breaks through all three limits by distributing data across multiple nodes. Each node handles a fraction of the total workload.


Key-Range Partitioning

How It Works

Assign a continuous range of keys to each partition, similar to volumes of an encyclopedia:

Partition 1: keys A-E
Partition 2: keys F-J
Partition 3: keys K-O
Partition 4: keys P-T
Partition 5: keys U-Z

The ranges are not necessarily evenly spaced -- they are chosen to distribute data evenly based on the actual key distribution.

Strengths
  • Efficient range queries: All keys in a range are on the same partition, so range scans are local and fast
  • Natural ordering: Data is stored in sorted order within each partition, supporting ORDER BY queries
  • Good for time-series: Partitioning by time range keeps recent data together
Weaknesses
  • Hotspots on sequential keys: If the partition key is a timestamp, all writes go to the partition for the current time period, creating a write hotspot
  • Uneven distribution: Key ranges that looked balanced at partition creation may become skewed as data grows
  • Manual or complex rebalancing: Range boundaries may need adjustment as data distribution changes
Avoiding Time-Series Hotspots

Instead of partitioning by timestamp alone, use a composite key:

Partition key: (sensor_id, date)

Sensor 1, 2024-01-01 -> Partition A
Sensor 2, 2024-01-01 -> Partition B
Sensor 1, 2024-01-02 -> Partition A
Sensor 3, 2024-01-01 -> Partition C

This distributes writes across partitions (different sensors go to different partitions) while preserving the ability to scan one sensor's data in time order.

Databases Using Key-Range Partitioning
  • HBase (row key ranges)
  • Bigtable (row key ranges)
  • MongoDB (range-based sharding option)

Hash Partitioning

How It Works

Apply a hash function to the partition key and assign hash ranges to partitions:

hash(key) mod N = partition number

hash("user_123") = 0x7A3F... -> Partition 3
hash("user_456") = 0x1B2C... -> Partition 1
hash("user_789") = 0xE4D1... -> Partition 5
Consistent Hashing

Standard hash mod N is problematic when adding or removing nodes because it reassigns most keys. Consistent hashing solves this by mapping both keys and nodes onto a ring:

Ring positions: 0 ... 2^32

Nodes:  A at position 1000, B at position 5000, C at position 9000
Key:    hash("user_123") = 3500 -> assigned to Node B (next node clockwise)

When a node is added or removed, only the keys between adjacent nodes are reassigned, minimizing data movement.

Virtual Nodes (Vnodes)

Each physical node is assigned multiple positions (virtual nodes) on the ring, typically 256 per node. This:

  • Distributes data more evenly (random positions may cluster otherwise)
  • Enables proportional assignment (a more powerful node gets more virtual nodes)
  • Smooths rebalancing (adding a node moves small chunks from many existing nodes)
Strengths
  • Even distribution: Hash functions distribute keys uniformly, avoiding hotspots from key distribution skew
  • Simple assignment: Given the key, you can compute the partition without a lookup table
Weaknesses
  • No range queries: Hash destroys sort order, so range scans require querying all partitions (scatter-gather)
  • Hot keys still possible: If one key receives disproportionate traffic (e.g., a celebrity's user ID), hashing doesn't help
Databases Using Hash Partitioning
  • Cassandra (default partitioner: Murmur3)
  • DynamoDB (hash of partition key)
  • Riak (consistent hashing)
  • MongoDB (hash-based sharding option)

Secondary Index Partitioning

When you need to query data by something other than the partition key, you need secondary indexes. There are two approaches to partitioning secondary indexes.

Local Secondary Indexes (Document-Partitioned)

Each partition maintains its own secondary index covering only the data in that partition:

Partition 1: primary data A-M, local index on "color"
Partition 2: primary data N-Z, local index on "color"

Query: SELECT * WHERE color = 'red'
  -> Must query BOTH partitions (scatter-gather)
  -> Each checks its local index
  -> Results are merged

Strengths:

  • Writes are local: updating the secondary index only affects one partition
  • Simple to maintain: each partition is self-contained

Weaknesses:

  • Reads require scatter-gather across all partitions
  • Latency is determined by the slowest partition (tail latency)

Used by: MongoDB, Cassandra, Elasticsearch, SolrCloud

Global Secondary Indexes (Term-Partitioned)

The secondary index is itself partitioned, but independently of the primary data:

Primary data: partitioned by user_id
Global index on "color": partitioned by color value
  Index partition 1: color A-M (all reds across all primary partitions)
  Index partition 2: color N-Z

Query: SELECT * WHERE color = 'red'
  -> Query index partition 1 only
  -> Get list of document IDs
  -> Fetch documents from their primary partitions

Strengths:

  • Reads are efficient: query only the relevant index partition
  • No scatter-gather for indexed queries

Weaknesses:

  • Writes require updating a remote partition (cross-partition write)
  • Index updates are often asynchronous, meaning the index may be stale
  • More complex distributed transaction requirements

Used by: DynamoDB (global secondary indexes), Amazon Aurora


Rebalancing Strategies

As data grows or nodes are added/removed, partitions must be rebalanced.

Strategy 1: Fixed Number of Partitions

Create many more partitions than nodes (e.g., 1000 partitions for 10 nodes). Each node hosts multiple partitions. When a node is added, some partitions move from existing nodes to the new node.

Before: 3 nodes, 12 partitions (4 per node)
After adding node 4: 4 nodes, 12 partitions (3 per node)
Move 1 partition from each existing node to the new node

Strengths: Simple, no re-partitioning needed, proportional load balancing Weaknesses: Must choose partition count upfront; too few means large partitions, too many means overhead

Used by: Elasticsearch, Riak, Couchbase, Voldemort

Strategy 2: Dynamic Partitioning

Start with one partition. When a partition grows beyond a threshold (e.g., 10GB), split it in half. When it shrinks below a threshold, merge it with a neighbor.

Strengths: Adapts to data size automatically; no upfront sizing decisions Weaknesses: Single partition at start means single-node bottleneck until first split; can cause split storms under rapid growth

Used by: HBase, RethinkDB, MongoDB (with key-range sharding)

Strategy 3: Proportional to Nodes

Keep a fixed number of partitions per node. When a node is added, it splits some existing partitions; when removed, its partitions are merged into others.

Strengths: Partition count grows with cluster size; each partition stays manageable Weaknesses: Splitting introduces brief unavailability for the affected partition

Used by: Cassandra (with vnodes)


Request Routing

How does a client know which node holds the partition for a given key?

Approach 1: Client-Side Routing

The client knows the partition assignment and connects directly to the correct node:

Client: hash("user_123") -> Partition 3 -> Node B
Client connects directly to Node B

Requires the client to maintain a copy of the partition map. Used by Cassandra drivers.

Approach 2: Routing Tier (Proxy)

A separate routing tier receives all requests and forwards them to the correct node:

Client -> Proxy -> determines partition -> forwards to correct Node

Used by: MongoDB (mongos router), Twemproxy (for Redis/Memcached)

Approach 3: Any-Node Contact

Client contacts any node; that node forwards the request if it doesn't own the partition:

Client -> Node A -> "Not my partition" -> forwards to Node B

Used by: Cassandra (coordinator pattern), CockroachDB

Service Discovery

All approaches need to know the current partition-to-node mapping. Options:

  • ZooKeeper/etcd: Centralized configuration service that tracks which partitions are on which nodes; nodes register themselves; routing layer watches for changes
  • Gossip protocol: Nodes gossip partition assignments to each other; eventually consistent but no central point of failure
  • DNS-based: Simple but slow to update; suitable only for coarse-grained routing

Handling Hotspots

Why Hotspots Occur

Even with perfect hash distribution, application-level access patterns create hotspots:

  • Celebrity problem: A single user or entity receives vastly more traffic than others
  • Temporal hotspots: Events cause sudden spikes on specific keys (product launch, breaking news)
  • Sequential keys: Auto-incrementing IDs or timestamps concentrate writes
Mitigation Strategies
Strategy How It Works Trade-off
Key splitting Append random suffix (0-9) to hot keys; read from all 10 sub-keys and merge 10x fan-out on reads; application complexity
Write buffering Buffer writes to hot keys in memory; flush periodically Eventual consistency; risk of data loss if buffer crashes
Caching layer Cache hot reads in Redis/Memcached in front of the database Stale data; cache invalidation complexity
Rate limiting Limit requests to hot keys per client Degrades user experience for hot content
Application-level sharding Route hot entities to dedicated, scaled infrastructure Operational complexity; special-case architecture
Detecting Hotspots

Monitor per-partition metrics:

  • Request rate per partition: Compare against average; alert on 10x deviation
  • Latency per partition: Hot partitions show higher p99 latency
  • CPU/IO utilization per node: Uneven utilization signals partition skew
  • Key-level access counting: Sample or log the most-accessed keys (most databases provide slow query logs or key-access statistics)
Automatic Hotspot Detection

Some systems detect and mitigate hotspots automatically:

  • DynamoDB Adaptive Capacity: Automatically isolates hot partitions onto dedicated throughput
  • Spanner: Splits hot partitions when load exceeds threshold
  • CockroachDB: Automatic range splitting and lease rebalancing based on load
1# Partitioning
2 
3Partitioning (also called sharding) divides a large dataset into smaller subsets called partitions, each stored on a different node. The goal is to spread data and query load evenly across machines, enabling horizontal scaling beyond the capacity of a single node.
4 
5## Why Partition?
6 
7A single database node has hard limits:
8- **Storage capacity:** A single disk or SSD has a maximum size
9- **Write throughput:** A single CPU can process a limited number of writes per second
10- **Read throughput:** Even with caching, a single node can serve a limited number of concurrent reads
11 
12Partitioning breaks through all three limits by distributing data across multiple nodes. Each node handles a fraction of the total workload.
13 
14---
15 
16## Key-Range Partitioning
17 
18### How It Works
19 
20Assign a continuous range of keys to each partition, similar to volumes of an encyclopedia:
21 
22```
23Partition 1: keys A-E
24Partition 2: keys F-J
25Partition 3: keys K-O
26Partition 4: keys P-T
27Partition 5: keys U-Z
28```
29 
30The ranges are not necessarily evenly spaced -- they are chosen to distribute data evenly based on the actual key distribution.
31 
32### Strengths
33 
34- **Efficient range queries:** All keys in a range are on the same partition, so range scans are local and fast
35- **Natural ordering:** Data is stored in sorted order within each partition, supporting ORDER BY queries
36- **Good for time-series:** Partitioning by time range keeps recent data together
37 
38### Weaknesses
39 
40- **Hotspots on sequential keys:** If the partition key is a timestamp, all writes go to the partition for the current time period, creating a write hotspot
41- **Uneven distribution:** Key ranges that looked balanced at partition creation may become skewed as data grows
42- **Manual or complex rebalancing:** Range boundaries may need adjustment as data distribution changes
43 
44### Avoiding Time-Series Hotspots
45 
46Instead of partitioning by timestamp alone, use a composite key:
47 
48```
49Partition key: (sensor_id, date)
50 
51Sensor 1, 2024-01-01 -> Partition A
52Sensor 2, 2024-01-01 -> Partition B
53Sensor 1, 2024-01-02 -> Partition A
54Sensor 3, 2024-01-01 -> Partition C
55```
56 
57This distributes writes across partitions (different sensors go to different partitions) while preserving the ability to scan one sensor's data in time order.
58 
59### Databases Using Key-Range Partitioning
60 
61- HBase (row key ranges)
62- Bigtable (row key ranges)
63- MongoDB (range-based sharding option)
64 
65---
66 
67## Hash Partitioning
68 
69### How It Works
70 
71Apply a hash function to the partition key and assign hash ranges to partitions:
72 
73```
74hash(key) mod N = partition number
75 
76hash("user_123") = 0x7A3F... -> Partition 3
77hash("user_456") = 0x1B2C... -> Partition 1
78hash("user_789") = 0xE4D1... -> Partition 5
79```
80 
81### Consistent Hashing
82 
83Standard `hash mod N` is problematic when adding or removing nodes because it reassigns most keys. Consistent hashing solves this by mapping both keys and nodes onto a ring:
84 
85```
86Ring positions: 0 ... 2^32
87 
88Nodes: A at position 1000, B at position 5000, C at position 9000
89Key: hash("user_123") = 3500 -> assigned to Node B (next node clockwise)
90```
91 
92When a node is added or removed, only the keys between adjacent nodes are reassigned, minimizing data movement.
93 
94### Virtual Nodes (Vnodes)
95 
96Each physical node is assigned multiple positions (virtual nodes) on the ring, typically 256 per node. This:
97- Distributes data more evenly (random positions may cluster otherwise)
98- Enables proportional assignment (a more powerful node gets more virtual nodes)
99- Smooths rebalancing (adding a node moves small chunks from many existing nodes)
100 
101### Strengths
102 
103- **Even distribution:** Hash functions distribute keys uniformly, avoiding hotspots from key distribution skew
104- **Simple assignment:** Given the key, you can compute the partition without a lookup table
105 
106### Weaknesses
107 
108- **No range queries:** Hash destroys sort order, so range scans require querying all partitions (scatter-gather)
109- **Hot keys still possible:** If one key receives disproportionate traffic (e.g., a celebrity's user ID), hashing doesn't help
110 
111### Databases Using Hash Partitioning
112 
113- Cassandra (default partitioner: Murmur3)
114- DynamoDB (hash of partition key)
115- Riak (consistent hashing)
116- MongoDB (hash-based sharding option)
117 
118---
119 
120## Secondary Index Partitioning
121 
122When you need to query data by something other than the partition key, you need secondary indexes. There are two approaches to partitioning secondary indexes.
123 
124### Local Secondary Indexes (Document-Partitioned)
125 
126Each partition maintains its own secondary index covering only the data in that partition:
127 
128```
129Partition 1: primary data A-M, local index on "color"
130Partition 2: primary data N-Z, local index on "color"
131 
132Query: SELECT * WHERE color = 'red'
133 -> Must query BOTH partitions (scatter-gather)
134 -> Each checks its local index
135 -> Results are merged
136```
137 
138**Strengths:**
139- Writes are local: updating the secondary index only affects one partition
140- Simple to maintain: each partition is self-contained
141 
142**Weaknesses:**
143- Reads require scatter-gather across all partitions
144- Latency is determined by the slowest partition (tail latency)
145 
146**Used by:** MongoDB, Cassandra, Elasticsearch, SolrCloud
147 
148### Global Secondary Indexes (Term-Partitioned)
149 
150The secondary index is itself partitioned, but independently of the primary data:
151 
152```
153Primary data: partitioned by user_id
154Global index on "color": partitioned by color value
155 Index partition 1: color A-M (all reds across all primary partitions)
156 Index partition 2: color N-Z
157 
158Query: SELECT * WHERE color = 'red'
159 -> Query index partition 1 only
160 -> Get list of document IDs
161 -> Fetch documents from their primary partitions
162```
163 
164**Strengths:**
165- Reads are efficient: query only the relevant index partition
166- No scatter-gather for indexed queries
167 
168**Weaknesses:**
169- Writes require updating a remote partition (cross-partition write)
170- Index updates are often asynchronous, meaning the index may be stale
171- More complex distributed transaction requirements
172 
173**Used by:** DynamoDB (global secondary indexes), Amazon Aurora
174 
175---
176 
177## Rebalancing Strategies
178 
179As data grows or nodes are added/removed, partitions must be rebalanced.
180 
181### Strategy 1: Fixed Number of Partitions
182 
183Create many more partitions than nodes (e.g., 1000 partitions for 10 nodes). Each node hosts multiple partitions. When a node is added, some partitions move from existing nodes to the new node.
184 
185```
186Before: 3 nodes, 12 partitions (4 per node)
187After adding node 4: 4 nodes, 12 partitions (3 per node)
188Move 1 partition from each existing node to the new node
189```
190 
191**Strengths:** Simple, no re-partitioning needed, proportional load balancing
192**Weaknesses:** Must choose partition count upfront; too few means large partitions, too many means overhead
193 
194**Used by:** Elasticsearch, Riak, Couchbase, Voldemort
195 
196### Strategy 2: Dynamic Partitioning
197 
198Start with one partition. When a partition grows beyond a threshold (e.g., 10GB), split it in half. When it shrinks below a threshold, merge it with a neighbor.
199 
200**Strengths:** Adapts to data size automatically; no upfront sizing decisions
201**Weaknesses:** Single partition at start means single-node bottleneck until first split; can cause split storms under rapid growth
202 
203**Used by:** HBase, RethinkDB, MongoDB (with key-range sharding)
204 
205### Strategy 3: Proportional to Nodes
206 
207Keep a fixed number of partitions per node. When a node is added, it splits some existing partitions; when removed, its partitions are merged into others.
208 
209**Strengths:** Partition count grows with cluster size; each partition stays manageable
210**Weaknesses:** Splitting introduces brief unavailability for the affected partition
211 
212**Used by:** Cassandra (with vnodes)
213 
214---
215 
216## Request Routing
217 
218How does a client know which node holds the partition for a given key?
219 
220### Approach 1: Client-Side Routing
221 
222The client knows the partition assignment and connects directly to the correct node:
223 
224```
225Client: hash("user_123") -> Partition 3 -> Node B
226Client connects directly to Node B
227```
228 
229Requires the client to maintain a copy of the partition map. Used by Cassandra drivers.
230 
231### Approach 2: Routing Tier (Proxy)
232 
233A separate routing tier receives all requests and forwards them to the correct node:
234 
235```
236Client -> Proxy -> determines partition -> forwards to correct Node
237```
238 
239Used by: MongoDB (mongos router), Twemproxy (for Redis/Memcached)
240 
241### Approach 3: Any-Node Contact
242 
243Client contacts any node; that node forwards the request if it doesn't own the partition:
244 
245```
246Client -> Node A -> "Not my partition" -> forwards to Node B
247```
248 
249Used by: Cassandra (coordinator pattern), CockroachDB
250 
251### Service Discovery
252 
253All approaches need to know the current partition-to-node mapping. Options:
254 
255- **ZooKeeper/etcd:** Centralized configuration service that tracks which partitions are on which nodes; nodes register themselves; routing layer watches for changes
256- **Gossip protocol:** Nodes gossip partition assignments to each other; eventually consistent but no central point of failure
257- **DNS-based:** Simple but slow to update; suitable only for coarse-grained routing
258 
259---
260 
261## Handling Hotspots
262 
263### Why Hotspots Occur
264 
265Even with perfect hash distribution, application-level access patterns create hotspots:
266 
267- **Celebrity problem:** A single user or entity receives vastly more traffic than others
268- **Temporal hotspots:** Events cause sudden spikes on specific keys (product launch, breaking news)
269- **Sequential keys:** Auto-incrementing IDs or timestamps concentrate writes
270 
271### Mitigation Strategies
272 
273| Strategy | How It Works | Trade-off |
274|----------|-------------|-----------|
275| **Key splitting** | Append random suffix (0-9) to hot keys; read from all 10 sub-keys and merge | 10x fan-out on reads; application complexity |
276| **Write buffering** | Buffer writes to hot keys in memory; flush periodically | Eventual consistency; risk of data loss if buffer crashes |
277| **Caching layer** | Cache hot reads in Redis/Memcached in front of the database | Stale data; cache invalidation complexity |
278| **Rate limiting** | Limit requests to hot keys per client | Degrades user experience for hot content |
279| **Application-level sharding** | Route hot entities to dedicated, scaled infrastructure | Operational complexity; special-case architecture |
280 
281### Detecting Hotspots
282 
283Monitor per-partition metrics:
284- **Request rate per partition:** Compare against average; alert on 10x deviation
285- **Latency per partition:** Hot partitions show higher p99 latency
286- **CPU/IO utilization per node:** Uneven utilization signals partition skew
287- **Key-level access counting:** Sample or log the most-accessed keys (most databases provide slow query logs or key-access statistics)
288 
289### Automatic Hotspot Detection
290 
291Some systems detect and mitigate hotspots automatically:
292- **DynamoDB Adaptive Capacity:** Automatically isolates hot partitions onto dedicated throughput
293- **Spanner:** Splits hot partitions when load exceeds threshold
294- **CockroachDB:** Automatic range splitting and lease rebalancing based on load
295 

Discussion