Section 3: Data Operations & Internals

Cassandra Data Distribution

Master the art of distributing petabytes across thousands of nodes! Deep dive into consistent hashing, token rings, replication strategies, and the mathematical elegance of distributed data placement.

๐Ÿ“– The Story: The Pizza Delivery Empire

Imagine MegaPizza - a global delivery company with 1 billion customers. How do you organize customer data so ANY order arrives in minutes?

โŒ BAD System: Single Central Database

The Problem: One massive server in New York stores all customer data

๐Ÿข NYC Server (OVERLOADED):
โ”œโ”€ Alice (LA) - 3000 miles away!
โ”œโ”€ Bob (Tokyo) - 7000 miles away!
โ”œโ”€ Charlie (London) - 3500 miles away!
โ””โ”€ ... 1 BILLION customers!

Problems:
โ€ข Latency: Tokyo user waits 200ms!
โ€ข Single point of failure: Server dies = business stops
โ€ข Scalability limit: Can't handle 1 billion!
โ€ข Hot spot: NYC server at 100% CPU

โœ… BRILLIANT System: Distributed Ring

The Solution: 100 regional hubs arranged in a circle (ring)

THE RING:

Hub-25(NYC) โ†’ Hub-26(Boston) โ†’ Hub-27(Miami)
โ†‘ โ†“
Hub-24(Seattle) Hub-28(Atlanta)
โ†‘ โ†“
Hub-23(LA) โ† Hub-1(Tokyo) โ† Hub-29(London)

Distribution Algorithm:
1. Customer name โ†’ Hash โ†’ Number (0-99)
2. Number determines which hub stores data
3. Data copies stored at 3 consecutive hubs

Examples:
Alice โ†’ hash("Alice") = 47 โ†’ Hub-47, Hub-48, Hub-49
Bob โ†’ hash("Bob") = 12 โ†’ Hub-12, Hub-13, Hub-14
Charlie โ†’ hash("Charlie") = 89 โ†’ Hub-89, Hub-90, Hub-91

Why It's GENIUS:

  • Even Distribution: Each hub stores ~10 million customers (perfect balance!)
  • Low Latency: Tokyo user gets data from Hub-1 (local!)
  • Fault Tolerant: Hub-47 crashes? Still have copies at Hub-48 & Hub-49!
  • Scalable: Add more hubs to ring โ†’ automatically rebalance!
  • No Hotspots: Similar names (Alice, Alicia) go to different hubs!

๐ŸŽฏ This is EXACTLY Cassandra's Data Distribution!

  • Pizza Hubs = Cassandra Nodes
  • Ring of Hubs = Token Ring
  • Hash Number = Token (64-bit integer)
  • Customer Name = Partition Key
  • 3 Consecutive Hubs = Replication Factor 3
CASSANDRA DATA DISTRIBUTION:

user_id="alice_2024" โ†’ Murmur3 โ†’ Token: -3,847,293,847,293
Token falls in Node 2's range โ†’ PRIMARY
RF=3 โ†’ Also stored on Node 3, Node 4 (clockwise)

Result: 3 copies across 3 nodes! โœ…

MegaPizza's hub system = Cassandra's token ring!
Same genius, distributed scale!

๐ŸŒ What is Data Distribution in Cassandra?

The mathematical framework that determines where every piece of data lives in the cluster.

Complete Definition

Data Distribution: The process by which Cassandra determines the physical location of data across nodes using consistent hashing, token assignments, and replication strategies to ensure even load distribution, fault tolerance, and linear scalability.

The Three Pillars:

  • Consistent Hashing: Maps data to tokens using hash functions
  • Token Ring: Organizes nodes in circular topology (-2^63 to 2^63-1)
  • Replication: Creates multiple copies across nodes for durability

Why It Matters:

  • Determines query latency (data locality)
  • Ensures even load across cluster
  • Enables horizontal scalability
  • Provides fault tolerance
  • Eliminates single points of failure
Data Distribution Complete Flow DATA ARRIVES INSERT INTO users (user_id='alice_2024') HASH Murmur3 Token: -3,847,293... TOKEN RING REPLICA PLACEMENT NODE 1 (Primary) Token: -4,000,000,000,000 โœ… Stores Data Owns this token range NODE 2 (Replica) Next clockwise โœ… Stores Copy RF=3, Replica 2 NODE 3 (Replica) Next clockwise โœ… Stores Copy RF=3, Replica 3 RESULT: 3 COPIES DISTRIBUTED! โ€ข Even distribution across nodes โ€ข Fault tolerant (can lose 1 node) โ€ข Fast reads (3 options to query)

โญ• Consistent Hashing: The Mathematical Foundation

Why Cassandra uses consistent hashing instead of traditional modulo-based distribution.

โŒ

Simple Hashing (BAD)

Traditional approach: node = hash(key) % N

N = 3 nodes
alice โ†’ hash=100 โ†’ 100%3=1 โ†’ Node 1
bob โ†’ hash=50 โ†’ 50%3=2 โ†’ Node 2

Add 4th node (N=4):
alice โ†’ 100%4=0 โ†’ Node 0 โŒ MOVED!
bob โ†’ 50%4=2 โ†’ Node 2 โœ… Same

Result: 75% of data needs to move!

Problem: Adding/removing nodes causes massive data migration!

โœ…

Consistent Hashing (GOOD)

Cassandra approach: Token ring with ranges

Token Ring (-2^63 to 2^63):
Node 1: -6T to -2T
Node 2: -2T to 2T
Node 3: 2T to 6T

Add Node 4:
Node 4 takes half of Node 1's range
Only ~25% of data moves!
alice stays on same node โœ…

Result: Minimal disruption!

Benefit: Adding nodes only affects immediate neighbors!

Expert Insight: Why Consistent Hashing is Revolutionary

The Problem It Solves:

In traditional distributed systems, adding/removing servers requires rehashing and relocating most data. With 1000 servers and 1PB data, adding one server means moving ~1TB of data!

Consistent Hashing Solution:

  • Data mapped to points on a ring (not to servers directly)
  • Servers claim ranges on the ring
  • Adding server: Only data in that range moves (~K/N where K=keys, N=nodes)
  • For 1000 nodes + 1PB data: Adding node moves only ~1GB!

Mathematical Guarantee:

With N nodes and K keys, adding/removing one node moves at most K/N keys. This is optimal - you can't do better!

๐Ÿ’ Token Ring: Cassandra's Circular Universe

Understanding the 64-bit token space and how nodes organize themselves.

Token Ring Visualization (6 Nodes) -9,223,372,036,854,775,808 0 9,223,372,036,854,775,807 Node 1 Token: -7.6T Node 2 Token: -4.6T Node 3 Token: -1.5T Node 4 Token: 1.5T Node 5 Token: 4.6T Node 6 Token: 7.6T TOKEN RING Clockwise Traversal โ†’ Each node owns range from its token to next node's token (clockwise)

Important: Token Ownership Rules

  • Each node has a token (its position on the ring)
  • Node owns all tokens from its token to next node's token (clockwise)
  • Ring wraps around (after max token comes min token)
  • Virtual nodes (vnodes): Each physical node gets multiple tokens
  • Default: 256 vnodes per node for better distribution

๐Ÿ–ฅ๏ธ Interactive Distribution Simulator

Experiment with data distribution in real-time!

Data Distribution Simulator
๐ŸŒ Distribution Simulator Ready!
Configure cluster settings and enter a partition key to see where data gets distributed.

Try these examples:
โ€ข user_12345
โ€ข order_abc123
โ€ข session_xyz789
โ€ข Click "Batch Test" to see 10 keys distributed!

๐ŸŒ Real-World Production Examples

How tech giants use data distribution at massive scale.

๐ŸŽฌ

Netflix

2,500 nodes, 1PB+ data, 30 million ops/sec

Keyspace: user_viewing
RF: 3 (per datacenter)
Datacenters: 3 (US-East, US-West, EU)

Example Distribution:
user_id='alice_123'
โ†’ Token: -2,847,239,847,234
โ†’ US-EAST: Node 47, 48, 49
โ†’ US-WEST: Node 128, 129, 130
โ†’ EU-WEST: Node 89, 90, 91

Total: 9 copies globally!

Why: Can lose entire datacenter and still serve!

๐ŸŽ

Apple iCloud

75,000+ nodes, 10PB+ data, 2 billion users

Keyspace: contacts
RF: 3 (per DC)
Vnodes: 256 per node

Distribution Strategy:
device_id='iPhone_xyz'
โ†’ Hash determines primary DC
โ†’ LOCAL_QUORUM for writes
โ†’ Sub-10ms latency globally

Load: ~10M keys per node
Balance: 99.8% uniformity!

Why: Perfect distribution = no hotspots!

๐Ÿš—

Uber

1,000+ nodes, 100TB+ data, real-time geolocation

Keyspace: driver_locations
RF: 3
Write: LOCAL_QUORUM

Distribution Pattern:
driver_id='driver_456'
โ†’ Token based on driver ID
โ†’ 3 replicas in local DC
โ†’ Updates every 4 seconds

Challenge: Time-based keys
Solution: Add driver_id to key!

Why: Geographic distribution + low latency!

๐Ÿ“Š Detailed Scenario Deep-Dives

Scenario 1: E-Commerce Platform - ShopGlobal.com

Company: ShopGlobal - Global marketplace with 500M users, 10M sellers

Scale: 100,000 orders/second peak, 50TB data, 200-node cluster

Table Design & Distribution Strategy

// User Orders Table
CREATE TABLE user_orders (
  user_id UUID,
  order_date DATE,
  order_id TIMEUUID,
  total_amount DECIMAL,
  status TEXT,
  items LIST<frozen<order_item>>,
  PRIMARY KEY ((user_id, order_date), order_id)
) WITH CLUSTERING ORDER BY (order_id DESC);

Replication Strategy:
RF = 3 per datacenter
Datacenters: US-EAST, US-WEST, EU-WEST, ASIA-PAC
Total replicas per order: 12 copies!

Real Distribution Example

Customer Order Flow:

Step 1: Order Placed
user_id = 'a1b2c3d4-e5f6-...'
order_date = '2024-12-26'
Partition Key: (user_id, order_date)

Step 2: Token Generation
Combined key serialized:
โ†’ Murmur3 hash
โ†’ Token: -4,892,374,928,374,928

Step 3: Ring Lookup (200 nodes)
Token range per node: ~92Q tokens
Token falls in Node 89's range

Step 4: Replica Placement
US-EAST (Primary DC):
  โ†’ Node 89 (primary)
  โ†’ Node 90 (replica 2)
  โ†’ Node 91 (replica 3)

US-WEST (Async replication):
  โ†’ Node 145, 146, 147

EU-WEST:
  โ†’ Node 67, 68, 69

ASIA-PAC:
  โ†’ Node 23, 24, 25

Performance Metrics

  • Write Latency: 3-5ms (LOCAL_QUORUM)
  • Read Latency: 1-2ms (LOCAL_ONE)
  • Data per Node: ~250GB (perfectly balanced)
  • Hotspot Prevention: Date in partition key spreads orders evenly
  • Fault Tolerance: Can lose 2 nodes per DC without data loss

Why This Design Works

  • Composite Partition Key: (user_id, order_date) ensures orders distributed across nodes AND time
  • No Time-Based Hotspots: Different dates = different nodes
  • Query Efficiency: "Get user's orders for date" = single partition read
  • Global Distribution: 4 DCs mean 12 copies, survive entire region failure
  • Linear Scalability: Add nodes โ†’ automatic rebalancing

Scenario 2: Social Media Platform - SocialVerse

Company: SocialVerse - Social network with 2B users, 100M DAU

Scale: 500,000 posts/second, 1PB data, 500-node cluster

Challenge: Celebrity Hotspot Problem

โŒ Initial Design (BROKEN):

CREATE TABLE user_timeline (
  user_id UUID PRIMARY KEY,
  posts LIST<frozen<post>>
);

Problem with Celebrity (100M followers):
All 100M follower timelines point to SAME celebrity user_id
โ†’ Same token
โ†’ Same 3 nodes (with RF=3)
โ†’ HOTSPOT! Nodes at 100% CPU while others idle!

Example:
celebrity_id='taylor_swift'
โ†’ Token: -5,123,456,789,012
โ†’ ALWAYS Node 234, 235, 236
โ†’ These 3 nodes handle ALL 100M requests!
โ†’ CLUSTER MELTDOWN! ๐Ÿ’ฅ

โœ… Fixed Design (BRILLIANT):

CREATE TABLE user_timeline_sharded (
  user_id UUID,
  shard_id INT, // 0-99
  post_id TIMEUUID,
  content TEXT,
  likes_count COUNTER,
  PRIMARY KEY ((user_id, shard_id), post_id)
) WITH CLUSTERING ORDER BY (post_id DESC);

Solution: Partition Key Sharding!
celebrity_id + shard_id = different tokens!

taylor_swift + shard_0 โ†’ Token A โ†’ Node 45, 46, 47
taylor_swift + shard_1 โ†’ Token B โ†’ Node 123, 124, 125
taylor_swift + shard_2 โ†’ Token C โ†’ Node 289, 290, 291
...
taylor_swift + shard_99 โ†’ Token Z โ†’ Node 456, 457, 458

Result: 100M requests spread across 300 nodes!
Load per node: ~333K requests (manageable!)
HOTSPOT ELIMINATED! โœ…

Application Logic

// Writing a post
shard_id = hash(post_id) % 100;
INSERT INTO user_timeline_sharded
(user_id, shard_id, post_id, content)
VALUES ('taylor_swift', shard_id, now(), 'Hello!');

// Reading timeline
// Query ALL shards in parallel
FOR shard_id IN 0..99:
  SELECT * FROM user_timeline_sharded
  WHERE user_id='taylor_swift'
  AND shard_id=shard_id
  LIMIT 100;

// Merge results client-side
// Sort by post_id (timestamp)
// Return top 100

Performance Impact

  • Before Sharding: 3 nodes at 100% CPU, 497 nodes idle
  • After Sharding: Load distributed across 300 nodes (~60% of cluster)
  • Write Latency: Reduced from 500ms to 5ms
  • Read Latency: 100 parallel queries = 10ms total (async)
  • Scalability: Can now handle 10x more celebrities!

Scenario 3: IoT Platform - SmartHome Inc.

Company: SmartHome - IoT platform with 50M devices, 100K new sensors/day

Scale: 1M sensor readings/second, 500TB timeseries data, 300-node cluster

Time-Series Data Distribution

// Sensor Data Table
CREATE TABLE sensor_readings (
  device_id UUID,
  date DATE,
  reading_time TIMESTAMP,
  temperature DOUBLE,
  humidity DOUBLE,
  battery_level INT,
  PRIMARY KEY ((device_id, date), reading_time)
) WITH CLUSTERING ORDER BY (reading_time DESC)
AND compaction = {
  'class': 'TimeWindowCompactionStrategy',
  'compaction_window_size': '1',
  'compaction_window_unit': 'DAYS'
};

Replication: RF=3, NetworkTopologyStrategy

Distribution Analysis

50M Devices ร— 365 Days = 18.25 BILLION Partitions!

Example: 3 Thermostats on Same Day

Device 1:
device_id='thermostat_a1b2'
date='2024-12-26'
โ†’ Token: -8,234,567,890,123
โ†’ Node 45, 46, 47

Device 2:
device_id='thermostat_c3d4'
date='2024-12-26'
โ†’ Token: 2,345,678,901,234
โ†’ Node 178, 179, 180

Device 3:
device_id='thermostat_e5f6'
date='2024-12-26'
โ†’ Token: 6,789,012,345,678
โ†’ Node 256, 257, 258

Perfect Distribution!
Same date, different devices โ†’ different nodes!
No time-based hotspots!

Query Patterns & Performance

Common Queries:

Query 1: Latest Reading (Single Partition)
SELECT * FROM sensor_readings
WHERE device_id='thermostat_a1b2'
AND date='2024-12-26'
LIMIT 1;

Performance: 1-2ms (single node)
Distribution: Node 45 serves directly

Query 2: Day's History (Single Partition)
SELECT * FROM sensor_readings
WHERE device_id='thermostat_a1b2'
AND date='2024-12-26';

Performance: 5-10ms (1 partition, ~17K rows)
Distribution: All data on Node 45, 46, 47

Query 3: Week's Data (7 Partitions)
SELECT * FROM sensor_readings
WHERE device_id='thermostat_a1b2'
AND date IN ('2024-12-20', ..., '2024-12-26');

Performance: 15-25ms (7 parallel queries)
Distribution: Could hit 7 different node groups!
Total rows: ~120K (perfect for cassandra)

Compaction Strategy Impact

  • TimeWindowCompactionStrategy: Perfect for time-series data
  • Window Size: 1 day = one SSTable per day per partition
  • Old Data: Drop entire SSTables (no tombstones!)
  • Retention: Delete data older than 90 days = just delete files
  • Performance: No compaction overhead during writes!

Capacity Planning

  • Data per Device per Day: ~100KB (1 reading/min ร— 1440 min ร— 70 bytes)
  • Total per Day: 50M devices ร— 100KB = 5TB/day
  • 90-Day Retention: 450TB raw data
  • With RF=3: 1.35PB total storage
  • Per Node (300 nodes): 4.5TB/node (comfortable!)
  • Distribution: 18.25B partitions รท 300 nodes = 60M partitions/node

Scenario 4: Financial Trading Platform - TradeFlow

Company: TradeFlow - High-frequency trading platform, 10M traders

Scale: 500K trades/second peak, 1TB/day audit logs, 100-node cluster

Multi-Table Distribution Strategy

// Table 1: Trade Executions (Query by trader)
CREATE TABLE trades_by_trader (
  trader_id UUID,
  trade_date DATE,
  trade_id TIMEUUID,
  symbol TEXT,
  quantity BIGINT,
  price DECIMAL,
  side TEXT, // BUY/SELL
  PRIMARY KEY ((trader_id, trade_date), trade_id)
);

// Table 2: Same data (Query by symbol)
CREATE TABLE trades_by_symbol (
  symbol TEXT,
  trade_date DATE,
  trade_id TIMEUUID,
  trader_id UUID,
  quantity BIGINT,
  price DECIMAL,
  side TEXT,
  PRIMARY KEY ((symbol, trade_date), trade_id)
);

Strategy: Duplicate data, different partition keys
โ†’ Different distributions!
โ†’ Optimized for different queries!

Distribution Comparison

Same Trade, Different Locations!

Trade Details:
trader_id='john_doe_123'
symbol='AAPL'
trade_date='2024-12-26'
trade_id=now()
price=150.25
quantity=1000

Distribution in trades_by_trader:
Partition Key: (john_doe_123, 2024-12-26)
โ†’ Token: -3,456,789,012,345
โ†’ Nodes: 34, 35, 36

Distribution in trades_by_symbol:
Partition Key: (AAPL, 2024-12-26)
โ†’ Token: 7,890,123,456,789
โ†’ Nodes: 78, 79, 80

Result: Same data, stored on DIFFERENT nodes!
This is OK! Optimizes for different query patterns!

Query Performance Analysis

  • Get Trader's Trades: Query trades_by_trader โ†’ Node 34 (1-2ms)
  • Get Symbol's Trades: Query trades_by_symbol โ†’ Node 78 (1-2ms)
  • Storage Overhead: 2x storage (worth it for query speed!)
  • Write Amplification: 2 writes per trade (batched, still fast)
  • Consistency: Both writes use QUORUM (strong consistency)

Hot Symbol Problem & Solution

Problem: AAPL has 1M trades/day

All AAPL trades for 2024-12-26 go to same partition!
โ†’ Same 3 nodes (78, 79, 80)
โ†’ Potential hotspot!

Solution: Time-based Micro-sharding

CREATE TABLE trades_by_symbol_hourly (
  symbol TEXT,
  trade_date DATE,
  hour INT, // 0-23
  trade_id TIMEUUID,
  ...
  PRIMARY KEY ((symbol, trade_date, hour), trade_id)
);

AAPL on 2024-12-26:
โ†’ (AAPL, 2024-12-26, 00) โ†’ Nodes 12, 13, 14
โ†’ (AAPL, 2024-12-26, 01) โ†’ Nodes 45, 46, 47
โ†’ (AAPL, 2024-12-26, 02) โ†’ Nodes 78, 79, 80
...
โ†’ (AAPL, 2024-12-26, 23) โ†’ Nodes 91, 92, 93

1M trades distributed across 72 nodes!
~14K trades per node (manageable!)

๐Ÿ’ผ Expert Interview Questions

Master these to demonstrate deep expertise!

1 Explain consistent hashing and why Cassandra uses it โ–ผ

Answer:

Consistent hashing maps both data and nodes to points on a circular hash space (ring), allowing minimal data movement when nodes are added or removed.

Why Cassandra Uses It:

  • Minimal Rebalancing: Adding node moves only K/N keys (K=total keys, N=nodes)
  • No Central Coordinator: Each node knows ring topology
  • Linear Scalability: Add nodes without downtime
  • Even Distribution: With vnodes, nearly perfect balance

Example:

Traditional: hash(key) % N
3 nodes โ†’ Add 4th node โ†’ 75% data moves!

Consistent Hashing:
Ring: -2^63 to 2^63-1
3 nodes โ†’ Add 4th node โ†’ 25% data moves!

With 100 nodes โ†’ Add 1 node โ†’ 1% moves!
2 How does replication factor affect data distribution? โ–ผ

Answer:

Replication Factor (RF) determines how many copies of each piece of data exist in the cluster. Data is replicated to RF consecutive nodes clockwise on the ring.

Distribution Rules:

  • RF=1: Data stored on 1 node only (no fault tolerance)
  • RF=3: Data stored on 3 consecutive nodes (industry standard)
  • RF=5: Data stored on 5 nodes (high availability scenarios)

Example with RF=3:

Key: user_id='alice'
Token: -3,847,293,847,293

Ring (6 nodes):
Node 1: -8T to -4T (PRIMARY - owns token)
Node 2: -4T to 0
Node 3: 0 to 4T
Node 4: 4T to 8T
...

With RF=3, replicas on:
โ†’ Node 1 (primary)
โ†’ Node 2 (next clockwise)
โ†’ Node 3 (next clockwise)

Result: Can lose 2 nodes and still have data!

Impact on Distribution:

  • Higher RF = More copies = Higher availability
  • Higher RF = More storage needed (RF=3 uses 3x space)
  • Higher RF = Slower writes (more nodes to coordinate)
  • Higher RF = Faster reads (more replicas to query)

๐ŸŽ“ Chapter Summary

You now understand Cassandra's data distribution!

Key Concepts:

  • Consistent hashing enables linear scalability
  • Token ring organizes 2^64 possible tokens
  • Replication provides fault tolerance
  • Vnodes ensure perfect distribution
  • Clockwise traversal for replica placement
Advertisement

Responsive Ad