NOTE
5.1 Distributed-System Partitioning: Data Splitting
Distribute keys across partitions/nodes as evenly as possible. 1. Choosing the Split Key 1.1. Primary-Key ID - For example, an auto-increment primary key. The advantage is even data distribution; the disadvantage is that queries by business fields may be slow because all partitions need to be read. 1.2. Business ID - Such as user ID or product ID.
This is a historical learning note and may contain outdated or incomplete understanding.
Distribute keys across partitions/nodes as evenly as possible.
1. Choosing the Split Key
1.1. Primary-Key ID
For example, an auto-increment primary key. The advantage is even data distribution; the disadvantage is that queries by business fields may be slow because all partitions need to be read.
1.2. Business ID
For example, user ID, product ID, and so on. The advantage is fast querying by business fields; the disadvantage is that data distribution may be uneven.
1.3. Examples
- MySQL is manually specified by the programmer and can use either a primary-key ID or a business ID.
- Redis itself is a key-value database, so it directly uses the key, which can be understood as a business key.
- Kafka uses a key, which can be understood as a business key.
- ZooKeeper does not use partitioning.
- Elasticsearch uses an automatically generated primary-key ID and can use a custom business ID.
2. Splitting Strategies
2.1. assigned partitioning
- Assigned partition
- Manually specify which partition a key is assigned to.
2.2. random partitioning
- Random partition
- Randomly assign a key to a partition.
2.3. range partitioning
- Sequential partition
- Specify a continuous key range (from minimum to maximum) for each partition.
- For example, for
user:1->user:N, assignuser:1-user:1000to partition1,user:1001-user:2000to partition2, and so on.
2.4. hash partitioning
2.4.1. hash%N
-
hash%N
- Calculate the hash value of the key, take the modulus by the total number of partitions, and place the data in the partition corresponding to the modulus result.
-
Advantages
- Simple and efficient.
-
Disadvantages
- If one machine goes down, (N-1)/N of the cache entries will miss.
-
Why is it (N-1)/N? For example, with 3 machines, hash values 1-6 are distributed across them as follows:
host 0: 3 6 host 1: 1 4 host 2: 2 5 -
If one machine goes down and only two remain, the modulus becomes 2 and the distribution becomes:
host 0: 2 4 6 host 1: 1 3 5 -
As shown, only two data items remain in the same location: 1 and 6. Four change location, which is 4/6 = 2/3 of the six data items.
-
- If one machine goes down, (N-1)/N of the cache entries will miss.
-
Use cases
- This approach is generally used for database sharding and partitioning, with doubling expansion to reduce the amount of data migration. Why not use consistent hashing? Consistent hashing needs the number of partitions to be determined, and if the number of partitions is fixed, later expansion is not possible.
- This approach is generally used for database sharding and partitioning, with doubling expansion to reduce the amount of data migration. Why not use consistent hashing? Consistent hashing needs the number of partitions to be determined, and if the number of partitions is fixed, later expansion is not possible.
2.4.2. consistent hash partitioning
- Consistent Hash
- Each partition is responsible for a range of hash values. Calculate the key’s hash value and determine which partition is responsible for the hash range containing that value.
2.4.2.1. hash ring partitioning
- Hash ring:
- Calculate the hash value of each server partition and place it on a circle containing values from 0 to 2^32-1.
- Calculate the hash value of the key and map it onto the same circle.
- Starting from the mapped position of the key, search clockwise and store it on the first server found.

- Advantages
- If one server goes down, the keys on that server move clockwise to the next server, and only 1/N will miss.
- Disadvantages
- Data skew: when there are too few servers, data is not evenly distributed around the circle and most keys may fall on the same partition, causing excessive load.
- Cannot solve the hot-key problem: solved through bounded consistent hashing.
2.4.2.1.1. hash ring + virtual node
- Calculate multiple hashes for each server partition and place them on the circle.
- The data-location algorithm remains unchanged.
- Advantage: solves the data-skew problem.
- Disadvantage: adding one partition causes more data migration. If all partitions are physical, adding one physical partition affects only one adjacent partition. With virtual partitions, the virtual partitions generated by adding one physical partition affect more partitions, but the affected ranges become smaller.
2.4.2.2. slot partitioning
- Hash slots:
- Redis Cluster does not use consistent hashing; instead it introduces the concept of hash slots.
- Redis Cluster has 16,384 hash slots for storing data, and each partition is responsible for some of the slots.
- A key uses
CRC16(key) % 16384to calculate which slot it belongs to, which then determines which partition it is assigned to.
- Hash ring vs slot partition: essentially, both add another layer to ensure that the position to which a key maps after hashing does not change when the number of partitions changes. Consistent Hash uses the hash directly for positioning, while Hash Slot uses a fixed modulus of 16,384.
2.5. random vs range vs hash
| random | range | hash | |
|---|---|---|---|
| Advantage | Distributes data evenly across all partitions | Range scans are very efficient and horizontal scaling is naturally supported | Key ranges are relatively evenly distributed |
| Disadvantage | Querying one data item requires parallel access to all partitions | Hot-data problems can occur; for example, newly added data may be concentrated in one partition and create a write bottleneck | Range scans are inefficient |
2.6. Examples
- MySQL is specified by the programmer and generally uses Hash%N.
- Redis uses slot partitioning in consistent hash.
- Kafka uses range by default, and can use Hash%N or an explicitly specified partition.
- ZooKeeper does not use partitioning.
- Elasticsearch uses Hash%N.


Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub