Files of Partitioning
wondelai/
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 | |
| 3 | 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. |
| 4 | |
| 5 | ## Why Partition? |
| 6 | |
| 7 | A 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 | |
| 12 | Partitioning 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 | |
| 20 | Assign a continuous range of keys to each partition, similar to volumes of an encyclopedia: |
| 21 | |
| 22 | |
| 23 | Partition 1: keys A-E |
| 24 | Partition 2: keys F-J |
| 25 | Partition 3: keys K-O |
| 26 | Partition 4: keys P-T |
| 27 | Partition 5: keys U-Z |
| 28 | |
| 29 | |
| 30 | The 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 | |
| 46 | Instead of partitioning by timestamp alone, use a composite key: |
| 47 | |
| 48 | |
| 49 | Partition key: (sensor_id, date) |
| 50 | |
| 51 | Sensor 1, 2024-01-01 -> Partition A |
| 52 | Sensor 2, 2024-01-01 -> Partition B |
| 53 | Sensor 1, 2024-01-02 -> Partition A |
| 54 | Sensor 3, 2024-01-01 -> Partition C |
| 55 | |
| 56 | |
| 57 | This 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 | |
| 71 | Apply a hash function to the partition key and assign hash ranges to partitions: |
| 72 | |
| 73 | |
| 74 | hash(key) mod N = partition number |
| 75 | |
| 76 | hash("user_123") = 0x7A3F... -> Partition 3 |
| 77 | hash("user_456") = 0x1B2C... -> Partition 1 |
| 78 | hash("user_789") = 0xE4D1... -> Partition 5 |
| 79 | |
| 80 | |
| 81 | ### Consistent Hashing |
| 82 | |
| 83 | 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: |
| 84 | |
| 85 | |
| 86 | Ring positions: 0 ... 2^32 |
| 87 | |
| 88 | Nodes: A at position 1000, B at position 5000, C at position 9000 |
| 89 | Key: hash("user_123") = 3500 -> assigned to Node B (next node clockwise) |
| 90 | |
| 91 | |
| 92 | When a node is added or removed, only the keys between adjacent nodes are reassigned, minimizing data movement. |
| 93 | |
| 94 | ### Virtual Nodes (Vnodes) |
| 95 | |
| 96 | Each 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 | |
| 122 | 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. |
| 123 | |
| 124 | ### Local Secondary Indexes (Document-Partitioned) |
| 125 | |
| 126 | Each partition maintains its own secondary index covering only the data in that partition: |
| 127 | |
| 128 | |
| 129 | Partition 1: primary data A-M, local index on "color" |
| 130 | Partition 2: primary data N-Z, local index on "color" |
| 131 | |
| 132 | Query: 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 | |
| 150 | The secondary index is itself partitioned, but independently of the primary data: |
| 151 | |
| 152 | |
| 153 | Primary data: partitioned by user_id |
| 154 | Global 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 | |
| 158 | Query: 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 | |
| 179 | As data grows or nodes are added/removed, partitions must be rebalanced. |
| 180 | |
| 181 | ### Strategy 1: Fixed Number of Partitions |
| 182 | |
| 183 | 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. |
| 184 | |
| 185 | |
| 186 | Before: 3 nodes, 12 partitions (4 per node) |
| 187 | After adding node 4: 4 nodes, 12 partitions (3 per node) |
| 188 | Move 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 | |
| 198 | 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. |
| 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 | |
| 207 | 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. |
| 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 | |
| 218 | How does a client know which node holds the partition for a given key? |
| 219 | |
| 220 | ### Approach 1: Client-Side Routing |
| 221 | |
| 222 | The client knows the partition assignment and connects directly to the correct node: |
| 223 | |
| 224 | |
| 225 | Client: hash("user_123") -> Partition 3 -> Node B |
| 226 | Client connects directly to Node B |
| 227 | |
| 228 | |
| 229 | Requires the client to maintain a copy of the partition map. Used by Cassandra drivers. |
| 230 | |
| 231 | ### Approach 2: Routing Tier (Proxy) |
| 232 | |
| 233 | A separate routing tier receives all requests and forwards them to the correct node: |
| 234 | |
| 235 | |
| 236 | Client -> Proxy -> determines partition -> forwards to correct Node |
| 237 | |
| 238 | |
| 239 | Used by: MongoDB (mongos router), Twemproxy (for Redis/Memcached) |
| 240 | |
| 241 | ### Approach 3: Any-Node Contact |
| 242 | |
| 243 | Client contacts any node; that node forwards the request if it doesn't own the partition: |
| 244 | |
| 245 | |
| 246 | Client -> Node A -> "Not my partition" -> forwards to Node B |
| 247 | |
| 248 | |
| 249 | Used by: Cassandra (coordinator pattern), CockroachDB |
| 250 | |
| 251 | ### Service Discovery |
| 252 | |
| 253 | All 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 | |
| 265 | Even 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 | |
| 283 | Monitor 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 | |
| 291 | Some 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
Browse more free Claude skills.