Cassandra Partitioner
Unlock the secret of perfect data distribution - master Murmur3, consistent hashing, and the algorithm that makes Cassandra scale to infinity!
📖 The Story: The Library's Organization Challenge
Meet The Grand Library with 10 million books spread across 100 buildings. The librarian, Alex, needed a system to decide which building stores each book.
❌ The Old System: Alphabetical (Simple but Broken)
How it worked: Books A-AA go to Building 1, AB-AZ to Building 2, B-BA to Building 3, etc.
Massive Problems:
- ⚠️ Uneven Distribution: Building 19 (S books) has 800K books (Shakespeare, Science, Self-help), Building 46 (X books) has only 15K books!
- 💥 Popular Letters Overload: Buildings for A, B, C, S are overwhelmed. Buildings for Q, X, Z are nearly empty.
- 🐌 Adding a Building: Insert new Building 50.5? Must renumber ALL buildings after 50 and move millions of books!
- 😫 Building Closes: Building 19 closes → 800K books must move to Building 20 (already full!) → Catastrophe!
Example Disaster: All "S" books in one building - when it closes for maintenance, students can't access Science, Shakespeare, Statistics, or any "S" book! System fails! 🔥
✅ The New System: Smart Hash Function (Partitioner!)
Brilliant Solution: Use a mathematical formula to convert book title to a number (0-999), then assign to building based on number!
How it works:
Hash("To Kill a Mockingbird") = 183 → Building 18
Hash("1984") = 91 → Building 9
Hash("Sapiens") = 456 → Building 45
Even though "Sapiens" starts with "S", it goes to Building 45!
Even though "1984" starts with "1", it goes to Building 9!
Distribution based on HASH, not alphabetical order!
Amazing Results:
- ⚖️ Perfect Balance: Each building has ~100,000 books (10M / 100 buildings)
- 🚀 Easy Expansion: Add Building 101? New books get hashed to 0-1000 range → automatically distributes to all 101 buildings
- 💪 Resilient: Building 74 closes? Only books with hash values 740-749 affected (~100K books), easily handled by neighbor buildings
- ⚡ Predictable Lookup: Hash any book title → know exactly which building, instantly!
Example Success: Building 74 closes. Books with hash 740-749 (100K books) automatically belong to Building 75. Simple, elegant, works! ✅
This is EXACTLY how Cassandra's Partitioner works!
Library = Cluster, Buildings = Nodes, Hash Function = Murmur3Partitioner
Perfect distribution, easy scaling, resilient to failures!
🔀 What is a Partitioner?
The master algorithm that determines where every piece of data lives in your Cassandra cluster.
Simple Definition
Partitioner: A hash function in Cassandra that takes your partition key and converts it into a token (a large number). This token determines which node in the cluster will store that data.
The Three-Step Process:
2. Partitioner hashes key: Hash("john_doe_123") = -3847293847293847289
3. Token determines node: Token -3847... belongs to Node 3 → Data stored on Node 3!
Analogy: Like a sorting hat in Harry Potter! It reads your name (partition key), thinks about it (hashes), and assigns you to a house (node). But instead of 4 houses, Cassandra has as many "houses" as you have nodes!
⚖️ Consistent Hashing: The Mathematical Magic
Understanding the algorithm that makes Cassandra scale infinitely without data chaos.
🎯 The Problem: Traditional Hashing Breaks at Scale
❌ Traditional Hash Function (node = hash(key) % N)
Scenario: You have 4 servers and 1 million user records.
user_001 → hash = 12345 → 12345 % 4 = 1 → Server 1
user_002 → hash = 67890 → 67890 % 4 = 2 → Server 2
user_003 → hash = 24680 → 24680 % 4 = 0 → Server 0
...
Result: Each server stores ~250,000 users ✅
The Disaster: Add 1 server (N=4 → N=5)
user_001 → 12345 % 5 = 0 → Server 0 (was Server 1!) 🔥
user_002 → 67890 % 5 = 0 → Server 0 (was Server 2!) 🔥
user_003 → 24680 % 5 = 0 → Server 0 (was Server 0!) ✅
...
Result: 80% of data (800K users) needs to MOVE!
Massive data migration, hours/days of downtime! 💥
✅ Consistent Hashing Solution
The Breakthrough: Hash both keys AND servers onto the same ring!
Key Principles:
- Fixed Hash Space: Use a huge ring (0 to 2^64-1)
- Server Positions: Hash each server name to position on ring
- Data Placement: Hash key, walk clockwise to find next server
- Add Server: Only data between new server and previous server moves!
Server 0 at position: 2,305,843,009,213,693,951
Server 1 at position: 4,611,686,018,427,387,903
Server 2 at position: 6,917,529,027,641,081,855
Server 3 at position: 9,223,372,036,854,775,807
user_001 hash = 3,000,000,000,000,000,000 → Server 1 ✅
Add Server 4 at position: 5,500,000,000,000,000,000
user_001 still hashes to 3,000... → Still Server 1! ✅
Only keys with hash 4,611... to 5,500... move to Server 4
Result: Only 20% of data (200K users) moves!
80% reduction in data movement! ⚡
How Consistent Hashing Works
Why Consistent Hashing is Brilliant
- Minimal Data Movement: Adding/removing node only affects 1/N of data (vs 80%+ in traditional hashing)
- Predictable: Same key always maps to same token (deterministic)
- Distributed: No central authority needed - each node independently calculates
- Scalable: Works equally well with 3 nodes or 3,000 nodes
- Self-Healing: Dead node's data automatically belongs to next clockwise node
🎯 Murmur3Partitioner: The Default Champion
Why Murmur3 became the gold standard for data distribution in Cassandra.
🏆 Meet Murmur3: The Perfect Hash Function
In 2012, Cassandra 1.2 introduced Murmur3Partitioner and it immediately became the default. Why? Because it's the Goldilocks of hash functions - not too slow, not too simple, just right!
What Makes Murmur3 Special?
- Blazing Fast: Processes ~3GB/sec on modern CPU (vs MD5: 0.5GB/sec)
- Perfect Distribution: Statistical tests show near-perfect uniformity
- Non-Cryptographic: Optimized for speed, not security (which we don't need!)
- Avalanche Effect: Change 1 bit in input → ~50% of output bits change
- 128-bit Hash: Produces massive output space (uses lower 64 bits as token)
How Murmur3 Generates Tokens
// Murmur3 Hash Algorithm (Simplified) 1. Take partition key: "user_12345" 2. Convert to bytes: [0x75, 0x73, 0x65, 0x72, 0x5f, 0x31, 0x32, 0x33, 0x34, 0x35] 3. Process in 128-bit blocks with mixing: - Mix with magic constants - Rotate bits (rotateLeft) - XOR operations for avalanche effect 4. Final output: -4826492847294729384 (64-bit signed token) Example Results: ================ Input: "user_1" → Token: -5234567890123456789 Input: "user_2" → Token: 7891234567890123456 Input: "user_3" → Token: -2345678901234567890 Notice: Similar inputs produce WILDLY different tokens! This ensures even distribution across token ring.
Speed Champion
Performance Metrics
MD5: 0.5 GB/sec
SHA256: 0.15 GB/sec
Winner: Murmur3 (6x faster!)
Distribution Quality
Statistical Tests
- Chi-squared test: Pass ✅
- Avalanche test: 50.1% ✅
- Collision rate: <0.0001% ✅
- Uniformity: 99.97% ✅
Configuration
cassandra.yaml
org.apache.cassandra
.dht.Murmur3Partitioner
# Default since Cassandra 1.2
📚 Types of Cassandra Partitioners
Understanding the evolution and choosing the right partitioner for your needs.
Murmur3Partitioner
✅ RECOMMENDED
Class:
Token Range:
-2^63 to 2^63-1
(Signed 64-bit)
When to Use:
- ✅ All new clusters
- ✅ Production systems
- ✅ High throughput
Pros:
- Fastest (3GB/s)
- Perfect distribution
- Industry standard
RandomPartitioner
⚠️ LEGACY
Class:
Token Range:
0 to 2^127-1
(Unsigned 128-bit)
When to Use:
- ⚠️ Legacy only
- ❌ Not for new
Cons:
- Slower MD5
- Larger tokens
- Legacy support
ByteOrderedPartitioner
❌ REMOVED
Status:
Removed in 3.0+
Why Failed:
- No hashing
- Hotspots
- Terrible balance
- Caused outages
Note:
Removed to prevent misuse
Cannot Change After Deployment!
Once deployed, partitioner is permanent! Changing requires building new cluster and migrating all data.
Best Practice: Always use Murmur3Partitioner!
⚙️ How Partitioner Works in Practice
Step-by-step journey from application to storage.
🔄 Complete Data Flow Example
Scenario: Storing User Profile
VALUES ('alice_2024', 'Alice', 'alice@example.com', 28);
Step 1: Extract Partition Key
Cassandra identifies user_id as partition key
Extracted: user_id = "alice_2024"
Step 2: Hash with Murmur3
Hash Process → MurmurHash3
Token: -3,847,293,847,293,847,289
Step 3: Find Owner Node
Node 1: -9,223... to -3,074...
Node 2: -3,074... to 3,074...
Node 3: 3,074... to 9,223...
Owner: Node 1
Step 4: Apply Replication (RF=3)
Replica 2: Node 2 (next)
Replica 3: Node 3 (next)
Write to all 3!
Step 5: Success!
Node 2: ✅ ACK (1ms)
Node 3: ⏱️ Writing...
QUORUM (2/3) → SUCCESS!
Total: ~5ms
🌍 Real-World Partitioner Examples
How industry giants leverage partitioners.
Photo Metadata
- Murmur3Partitioner
- Key: photo_id (UUID)
- 150+ nodes
- Petabytes of data
UUID keys = Perfect distribution across nodes!
Discord
Messages
- Murmur3Partitioner
- Key: (channel_id, bucket)
- 177 nodes
- Trillions of messages
Bucketing prevents hot channels!
Netflix
Viewing History
- Murmur3Partitioner
- Key: user_id
- 2,500+ nodes
- 300M+ users
120K users per node = Perfect balance!
Production Best Practices
- Always Murmur3 - Industry standard
- High-cardinality keys - user_id, UUID work great
- Avoid low-cardinality - country_code creates hotspots
- Composite keys - (user_id, bucket) for hot users
- Monitor distribution - Use nodetool status
💼 Interview Questions & Answers
Master these 30 essential partitioner questions!
Answer:
A partitioner is a hash function that determines data distribution across nodes. It converts partition keys into tokens, which map to specific nodes.
Why Critical:
- Load Balance: Even distribution prevents hotspots
- Scalability: Enables easy node addition/removal
- Performance: Efficient query routing
- Availability: Influences failure impact
Answer:
Consistent Hashing: Algorithm where data keys and servers hash onto the same circular ring. Data stored on first server encountered clockwise.
Benefits:
- Minimal disruption (1/N data moves vs 80%+)
- No resharding needed
- Predictable and deterministic
- Decentralized calculation
Answer:
Murmur3 became default in Cassandra 1.2 (2012) because:
- Performance: 3GB/s (6x faster than MD5)
- Distribution: 99.97%+ uniformity
- Non-cryptographic: Optimized for speed
- Avalanche effect: 1-bit change → 50% output changes
- Compatible: Works with vnodes
Answer:
ByteOrderedPartitioner used raw bytes as tokens (no hashing), creating catastrophic problems:
- Sequential keys: user_001, user_002 → same node
- Time-based disaster: All "today" data → one node
- Alphabet bias: "S" books overwhelm, "Q" books empty
- No vnodes: Couldn't work with virtual nodes
Removed: Too dangerous, caused production outages. Deprecated in 2.x, removed in 3.0+
Answer:
Method 1: nodetool getendpoints
10.0.0.3
10.0.0.5
10.0.0.2
Method 2: CQL Tracing
SELECT * FROM users WHERE user_id = 'alice_2024';
Answer:
NO! The partitioner is essentially permanent once a cluster is deployed. This is one of Cassandra's most critical design decisions.
Why It Can't Change:
- Token Assignment: All existing data mapped to specific tokens
- Data Location: Changing partitioner changes token → node mapping
- Complete Chaos: Would require moving ALL data to new locations
If You MUST Change (Extreme Case):
2. Use dual-write from application (write to both)
3. Use sstableloader to backfill old data to new cluster
4. Verify data integrity
5. Switch reads to new cluster
6. Decommission old cluster
Time required: Weeks to months
Risk: Very high
Cost: 2x infrastructure during migration
Best Practice: Choose wisely at cluster creation! Always use Murmur3Partitioner for new deployments. This decision is forever.
Answer:
Murmur3Partitioner uses the range: -2^63 to 2^63-1
In Numbers:
Maximum: +9,223,372,036,854,775,807
Total Space: 18,446,744,073,709,551,616 possible tokens
(That's 18.4 quintillion!)
Why This Range?
- 64-bit signed long: Standard Java/computer data type
- Efficient: Fits in single CPU register
- Fast operations: Native CPU arithmetic
- Large enough: 18 quintillion tokens = impossible to fill
- Balanced: Equal negative and positive space
Comparison with RandomPartitioner:
| Murmur3: | -2^63 to 2^63-1 (64-bit) |
| Random: | 0 to 2^127-1 (128-bit) |
| Murmur3's smaller range is more efficient while still providing massive space | |
Answer:
With composite partition keys, the partitioner hashes ALL components together as a single unit.
Example Table:
user_id TEXT,
event_date DATE,
event_id TIMEUUID,
event_type TEXT,
PRIMARY KEY ((user_id, event_date), event_id)
);
-- Composite partition key: (user_id, event_date)
Hashing Process:
"alice" + SEPARATOR + "2024-01-15"
2. Hash the combined value:
Murmur3("alice\x002024-01-15")
3. Result: Single token
Token: -5,234,567,890,123,456,789
Key Points:
- Single Token: Entire composite key produces ONE token
- Order Matters: (alice, 2024-01-15) ≠ (2024-01-15, alice)
- Same Node: All rows with same composite key on same node
- Distribution: Different composites get different tokens
Example Distribution:
(alice, 2024-01-16) → Token: 3,456... → Node 2
(bob, 2024-01-15) → Token: -8,901... → Node 3
Different dates for same user = Different nodes!
This spreads load even for hot users.
Answer:
Hash collisions are theoretically possible but practically impossible with Murmur3's 64-bit token space.
Math:
Collision probability (Birthday paradox):
- After 1 billion keys: ~0.00003% chance
- After 1 trillion keys: ~0.03% chance
Reality: Most clusters have < 100 billion keys
Collision probability: Effectively ZERO
If Collision Occurred (Theoretical):
- Same Node: Both keys map to same node (by design)
- Different Rows: They're still separate rows (partition key distinguishes)
- No Data Loss: Cassandra stores full partition key, not just token
- Performance: Minimal impact (slightly more data on one node)
Real-World:
In 10+ years of Cassandra usage across billions of deployments:
- Zero reported cases of hash collisions affecting operations
- The 64-bit space is astronomically large
- Not a concern in practice
Answer:
The partitioner determines the primary replica, then RF determines how many additional replicas to create by walking clockwise around the ring.
Process with RF=3:
2. Token maps to Node A (PRIMARY REPLICA)
3. Walk clockwise from Node A:
- Next node = Node B (REPLICA 2)
- Next node = Node C (REPLICA 3)
4. Data written to all 3 nodes
Example:
user_id="alice" → Token: -3,847,293,847,293,847,289
Token falls in Node 2's range (PRIMARY)
RF=3 requires 2 more replicas
Replicas:
1. Node 2 (primary - owns the token)
2. Node 3 (next clockwise)
3. Node 4 (next clockwise)
Query for alice can hit ANY of these 3 nodes!
With Virtual Nodes:
- Primary vnode determined by token
- Next RF-1 vnodes selected clockwise
- Vnodes can belong to different physical nodes
- Result: Better replica distribution
Network Topology Snitch: In multi-DC deployments, snitch ensures replicas spread across racks/DCs, but partitioner still determines the starting point (primary replica).
Answer:
Partition Key: The actual data value(s) you define (e.g., "alice_2024")
Token: The numeric hash result used for placement (e.g., -3,847,293,847,293,847,289)
Key Differences:
| Aspect | Partition Key | Token |
|---|---|---|
| Type | String, UUID, etc. | 64-bit integer |
| Purpose | Identify data | Placement decision |
| User-visible | Yes (in queries) | No (internal) |
| Storage | Stored in SSTable | Used for routing |
Example:
SELECT * FROM users WHERE user_id = 'alice_2024';
Behind the scenes:
1. Partition Key: 'alice_2024' (what you see)
2. Token: Murmur3('alice_2024') = -3,847,293... (internal)
3. Token maps to Node 2
4. Query routed to Node 2
5. Node 2 finds row with user_id='alice_2024'
Why Two Concepts?
- Abstraction: Users work with meaningful keys (user_id)
- Efficiency: System works with numeric tokens (fast comparison)
- Distribution: Tokens ensure even spread regardless of key values
Answer:
Method 1: nodetool status (Quick Check)
Datacenter: datacenter1
Status=Up/Down |/ State=Normal/Leaving/Joining
-- Address Load Tokens Owns Host ID
UN 10.0.0.1 1.2 TB 256 33.4% abc123...
UN 10.0.0.2 1.18 TB 256 33.1% def456...
UN 10.0.0.3 1.21 TB 256 33.5% ghi789...
✅ GOOD: All nodes own ~33.3% (3 nodes)
✅ GOOD: Load within 5% variance (1.18-1.21 TB)
❌ BAD: If one node had 45% and another 20%
Method 2: nodetool tablestats (Per-Table)
Shows per-table statistics:
- Space used
- Number of partitions
- Average partition size
Run on each node, compare values
What to Look For:
- Load Variance < 5%: Excellent distribution
- Load Variance 5-10%: Acceptable, monitor
- Load Variance > 10%: Investigation needed!
- Owns Percentage: Should be ~100/N for N nodes
Common Issues:
- Hot Partition: One key has massive data
- Poor Key Choice: Low cardinality partition key
- Data Skew: Some users much more active than others
- Sequential Keys: If using time-based keys incorrectly
Answer:
The avalanche effect means that changing even a single bit in the input causes approximately 50% of the output bits to flip. This is crucial for even data distribution.
Example:
Token: -5,234,567,890,123,456,789
Input: "user_12346" (one char different!)
Token: 7,891,234,567,890,123,456
Result: Completely different tokens!
Even though inputs are nearly identical
Why This Matters:
- Sequential Keys: user_001, user_002, user_003 get scattered across nodes
- Similar Data: alice_2024-01-01, alice_2024-01-02 go to different nodes
- No Clustering: Prevents sequential data from clustering on one node
- Even Distribution: Guarantees uniform spread regardless of key patterns
Without Avalanche Effect:
Bad Hash Function Example:
user_002 → Hash: 1001
user_003 → Hash: 1002
All similar values → Same node!
Creates hotspot! 🔥
Statistical Test: Murmur3 achieves 50.1% bit flipping (nearly perfect). This is why it's trusted for production at massive scale.
Answer:
Benchmark Comparison:
| Partitioner | Speed | CPU Impact | Latency |
|---|---|---|---|
| Murmur3 | 3 GB/s | Very Low | ~0.01ms |
| Random (MD5) | 0.5 GB/s | Medium | ~0.05ms |
Real-World Impact at Scale:
With Murmur3:
- Hash time: 0.01ms × 1M = 10 seconds CPU/sec
- CPU usage: ~1-2% per core
- Negligible impact ✅
With RandomPartitioner (MD5):
- Hash time: 0.05ms × 1M = 50 seconds CPU/sec
- CPU usage: ~5-8% per core
- Noticeable impact on high-load systems ⚠️
Why Speed Matters:
- Every Write: Partitioner called for every single write operation
- Every Read: Token calculation needed to route queries
- At Scale: Millions of operations/sec = CPU becomes bottleneck
- Latency: Faster hash = lower overall query latency
Netflix Example: By switching from RandomPartitioner to Murmur3, they reduced CPU usage by 3-4% across entire fleet, saving hundreds of thousands in infrastructure costs annually.
Answer:
The partitioner directly impacts query routing and coordinator efficiency.
Query Process:
2. Coordinator hashes partition key → Gets token
3. Coordinator looks up which nodes own that token
4. Coordinator routes query to correct replicas
5. Coordinator aggregates results
6. Returns to client
Token-Aware Drivers:
Modern drivers use partitioner information to connect DIRECTLY to the node owning the data:
Smart Routing:
Query for user_id='alice'
Driver calculates: Murmur3('alice') = Token X
Driver connects DIRECTLY to node owning Token X
Result: Zero-hop routing! ⚡
Latency reduced by 30-50%!
Performance Impact:
- Without Token-Aware: Query → Random node → Route to owner → Extra network hop
- With Token-Aware: Query → Direct to owner → No extra hop
- Latency Savings: 1-2ms per query (significant at scale)
- Load Distribution: No coordinator hotspots
Configuration:
Cluster cluster = Cluster.builder()
.addContactPoint("10.0.0.1")
.withLoadBalancingPolicy(
new TokenAwarePolicy(
DCAwareRoundRobinPolicy.builder().build()
)
)
.build();