Section 3: Data Distribution

Partitioning Basics

Master data distribution! Learn partition keys, consistent hashing, and how Cassandra distributes billions of rows with animated visual diagrams!

📖 Netflix's Partition Design: 500+ Million Users

Netflix stores viewing history for 500+ million users - billions of records. Challenge: How to distribute data evenly across 300 nodes? Solution: Perfect partition key design using user_id. Result: Every node stores ~1.67 million users, perfectly balanced load, sub-10ms queries at any scale!

🎯 What is Partitioning?

Partitioning is how Cassandra splits data across multiple nodes for horizontal scalability and parallel processing.

Partitioning: From Table to Nodes (Animated) Users Table (1 Billion Rows) user_id=1 | name='Alice' | email='alice@example.com' user_id=2 | name='Bob' | email='bob@example.com' user_id=3 | name='Charlie' | email='charlie@example.com' Hash Function Node 1 IP: 10.1.0.1 user_id=1 user_id=4 user_id=7 333M rows Node 2 IP: 10.1.0.2 user_id=2 user_id=5 user_id=8 333M rows Node 3 IP: 10.1.0.3 user_id=3 user_id=6 user_id=9 334M rows

Key Concepts

  • Partition Key: Determines which node stores the data (e.g., user_id)
  • Hash Function: Converts partition key to token (-2^63 to 2^63-1)
  • Token Ownership: Each node owns a range of tokens
  • Even Distribution: Hash function ensures balanced data across nodes
  • No Hot Spots: Data automatically spread evenly

🔑 Primary Key Design

The PRIMARY KEY has two components: partition key and clustering columns. This determines data distribution and query patterns!

Primary Key Anatomy (Animated) CREATE TABLE user_events PARTITION KEY user_id Determines node CLUSTERING KEY event_time Orders within partition REGULAR COLUMNS event_type, data Stored with partition PRIMARY KEY ((user_id), event_time) Common Primary Key Patterns Simple Partition Key PRIMARY KEY (user_id) • One row per partition • No ordering • Example: User profiles Composite Partition Key PRIMARY KEY ((user_id, date)) • Both determine node • Splits large partitions • Example: Daily metrics With Clustering Columns PRIMARY KEY ((user_id), timestamp) • Multiple rows per partition • Sorted by timestamp • Example: Time-series data Multiple Clustering Keys PRIMARY KEY ((sensor_id), date, hour) • Hierarchical ordering • Date first, then hour • Example: IoT sensor data Choose based on: Query patterns + Data size + Distribution needs Golden Rule: Partition key determines data location Clustering key determines data ordering within partition

Real-World Examples

# Example 1: User Profiles (Simple)
CREATE TABLE users (
    user_id uuid PRIMARY KEY,
    name text,
    email text,
    created_at timestamp
);
-- Partition key: user_id
-- Each user = 1 partition = 1 row
-- Perfect for lookups: SELECT * FROM users WHERE user_id = ?

# Example 2: User Activity Timeline (With Clustering)
CREATE TABLE user_activity (
    user_id uuid,
    activity_time timestamp,
    activity_type text,
    details text,
    PRIMARY KEY ((user_id), activity_time)
) WITH CLUSTERING ORDER BY (activity_time DESC);
-- Partition key: user_id
-- Clustering key: activity_time (sorted newest first)
-- Query: SELECT * FROM user_activity WHERE user_id = ? LIMIT 10
-- Returns: Last 10 activities for user

# Example 3: Time-Series Sensor Data (Composite Partition Key)
CREATE TABLE sensor_data (
    sensor_id uuid,
    date text,          -- '2024-12-26'
    hour int,
    minute int,
    temperature decimal,
    humidity decimal,
    PRIMARY KEY ((sensor_id, date), hour, minute)
);
-- Partition key: sensor_id + date (one partition per sensor per day)
-- Clustering keys: hour, minute (sorted chronologically)
-- Avoids huge partitions (splits by day)
-- Query: SELECT * FROM sensor_data 
--        WHERE sensor_id = ? AND date = '2024-12-26'

# Example 4: Social Media Posts (With Multiple Clustering)
CREATE TABLE user_posts (
    user_id uuid,
    post_date date,
    post_time timestamp,
    post_id uuid,
    content text,
    likes counter,
    PRIMARY KEY ((user_id, post_date), post_time, post_id)
) WITH CLUSTERING ORDER BY (post_time DESC, post_id ASC);
-- Partition key: user_id + post_date (one partition per user per day)
-- Clustering: post_time DESC (newest first), then post_id
-- Efficient queries: Get today's posts for user
-- Bounded partition size

⚙️ Consistent Hashing Function

Cassandra uses Murmur3 hash function to convert partition keys into tokens, ensuring even distribution!

Hash Function: Key → Token (Animated) Partition Keys user_id=1 user_id=2 user_id=3 user_id=4 user_id=5 Any value... String, UUID, Int... Murmur3 Hash Function • Deterministic • Uniform distribution • Fast computation • Range: -2^63 to 2^63-1 Tokens (64-bit) -8234719023847192 2847192038471920 -3847192084719201 7192084719201923 -9201923847192038 Evenly distributed across token range Hash Function Properties ✅ Same input always produces same token (deterministic) ✅ Different inputs produce evenly distributed tokens (uniform) ✅ Small change in input creates completely different token (avalanche effect)

Live Console: Viewing Tokens

# Check token for a specific partition key
$ cqlsh -e "SELECT token(user_id), user_id, name FROM users WHERE user_id = 12345"

 system.token(user_id) | user_id | name
-----------------------+---------+-------
  -8234719023847192847 |   12345 | Alice

# The token determines which node owns this data!

# View token ranges owned by each node
$ nodetool ring mykeyspace

Datacenter: US-East
Address      Rack  Status  State   Load      Owns   Token
10.1.0.1     rack1 Up      Normal  247.5 GB  33.3%  -9223372036854775808
10.1.0.2     rack1 Up      Normal  238.2 GB  33.2%  -3074457345618258603
10.1.0.3     rack2 Up      Normal  251.8 GB  33.5%   3074457345618258602
                                            ↑
                              Each node owns 1/3 of token range

# Token -8234719023847192847 falls between:
# -9223372036854775808 and -3074457345618258603
# Therefore user_id=12345 is stored on Node 10.1.0.1

# Get token for any value
$ cqlsh
cqlsh> SELECT token('alice@example.com');

 system.token('alice@example.com')
-----------------------------------
          2847192038471920384

cqlsh> SELECT token('bob@example.com');

 system.token('bob@example.com')
---------------------------------
        -3847192084719201923

# Different inputs → Different tokens → Different nodes!

Why Murmur3?

  • Speed: Extremely fast hash computation (~100ns)
  • Uniformity: Excellent distribution properties
  • Simplicity: Non-cryptographic (don't need security)
  • Avalanche Effect: Small input change → completely different output
  • Industry Standard: Used by Redis, Memcached, Hadoop

📊 Data Distribution Across Nodes

How Cassandra uses tokens to distribute data evenly across the cluster!

Token Ring: 3-Node Cluster (Animated) Token: 3074457345618258602 -9223372036854775808 -3074457345618258603 Node 3 10.1.0.3 Owns tokens: -3074...602 to 3074...602 33.3% of data Node 1 10.1.0.1 Owns tokens: 3074...602 to -9223...808 33.3% of data Node 2 10.1.0.2 Owns tokens: -9223...808 to -3074...603 33.4% of data user_id=1 user_id=2

Real Scenario: Spotify's 100-Node Cluster

Challenge: Store playlists for 500M users across 100 nodes

  • Total Token Range: -2^63 to 2^63-1 (18.4 quintillion values)
  • Per Node: Each owns ~1% of token range
  • Data Per Node: ~5 million users per node (500M / 100)
  • Distribution: Hash function ensures near-perfect 1% on each
  • Queries: token(user_id) determines node → single-node lookup!
# Example token distribution across 100 nodes
$ nodetool status | head -10

Address      Load      Owns   Token
10.1.0.1     245 GB    1.02%  -9223372036854775808
10.1.0.2     238 GB    0.98%  -9039503846923745923
10.1.0.3     251 GB    1.04%  -8855635656992716038
10.1.0.4     247 GB    1.01%  -8671767467061686153
...
10.1.0.100   243 GB    0.99%   9039503846923745922

# Perfect distribution! Each node ~1% of data (247GB avg)
# Total cluster: 24.7TB across 100 nodes

Console: Analyzing Data Distribution

# Check ownership percentages
$ nodetool status

Datacenter: Production
Address      Load      Owns     Token                    
10.1.0.1     247.5 GB  33.28%   -9223372036854775808
10.1.0.2     238.2 GB  33.24%   -3074457345618258603
10.1.0.3     251.8 GB  33.48%    3074457345618258602
            ↑ Good! Nearly equal distribution

# Check load per table
$ nodetool tablestats mykeyspace.users | grep "Space used"
Space used (total): 247534892347
Space used (live): 247234892347
Space used (SSTable): 247134892347

# Find which node owns a specific key
$ nodetool getendpoints mykeyspace users 12345
10.1.0.1
10.1.0.2  # If RF=2
10.1.0.3  # If RF=3

# Verify even distribution with token distribution
$ nodetool describering mykeyspace | grep "Range"
        Range(3074457345618258602, -9223372036854775808]
        Range(-3074457345618258603, 3074457345618258602]
        Range(-9223372036854775808, -3074457345618258603]
        ↑ Three equal ranges for 3-node cluster

🔥 Hot Partitions Problem

When one partition becomes too large or too frequently accessed, it creates a performance bottleneck!

✅

Good Partition

  • Size: < 100MB
  • Rows: < 100,000
  • Even access pattern
  • Fast queries (< 10ms)
  • Fits in memory
❌

Hot Partition

  • Size: > 1GB
  • Rows: > 1 million
  • Heavy traffic
  • Slow queries (> 100ms)
  • Causes node issues

Real Disaster: Twitter's Celebrity Problem

Problem: Twitter's follower table using celebrity_id as partition key

# BAD DESIGN:
CREATE TABLE followers (
    celebrity_id uuid,
    follower_id uuid,
    followed_at timestamp,
    PRIMARY KEY ((celebrity_id), follower_id)
);

# Barack Obama has 130+ million followers
# This creates ONE MASSIVE PARTITION:

Partition: celebrity_id = obama_uuid
├── follower_1 (130 million rows!)
├── follower_2
├── ...
└── follower_130000000

Partition Size: 15 GB (way too big!)
Query Time: 5+ seconds to scan
Node Impact: Causes OOM errors, GC pauses

# Query that kills the cluster:
SELECT * FROM followers 
WHERE celebrity_id = obama_uuid;
-- This reads 15GB from ONE node!
-- Node runs out of memory
-- Cluster destabilized

✅ Solution: Partition by follower_id instead!

# GOOD DESIGN:
CREATE TABLE user_following (
    follower_id uuid,
    celebrity_id uuid,
    followed_at timestamp,
    PRIMARY KEY ((follower_id), celebrity_id)
);

# Now each user has small partition
Partition: follower_id = alice_uuid
├── follows celebrity_1
├── follows celebrity_2
├── follows celebrity_3 (maybe 50 rows total)

Partition Size: < 100 KB (perfect!)
Query Time: < 1ms
Distribution: Even across all nodes!

# Query pattern changes:
SELECT * FROM user_following 
WHERE follower_id = alice_uuid;
-- Reads < 100KB from ONE partition
-- Lightning fast!
-- Scales to billions of users!

Detecting Hot Partitions

# Find large partitions
$ nodetool cfstats mykeyspace.users | grep -A 5 "SSTable count"

SSTable count: 47
Space used (live): 247534892347
Space used (total): 247634892347
Space used by snapshots (total): 0
Off heap memory used (total): 142357234
Bloom filter false positives: 0
Bloom filter false ratio: 0.00000
Compacted partition minimum bytes: 152
Compacted partition maximum bytes: 15847234567  ← HUGE!
                                     ↑ This partition is 15GB!
Compacted partition mean bytes: 247534

# Find partitions with most rows
$ nodetool tablehistograms mykeyspace.users

users histograms
Percentile  SSTables   Write Latency  Read Latency  Partition Size  Cell Count
           (micros)      (micros)         (micros)      (bytes)
50%            1.00          12.00         89.00         1024           10
75%            2.00          18.00        125.00         4096           25
95%            4.00          45.00        456.00        16384          100
98%            6.00          78.00        892.00        65536          500
99%            8.00         124.00       1847.00      524288         2000
Min            1.00           8.00         12.00          256            1
Max           47.00        4892.00      45892.00   15847234567   130000000
                                                           ↑              ↑
                                                    15GB partition  130M rows!

# Monitor read/write hotspots
$ nodetool proxyhistograms

Percentile      Read Latency     Write Latency     Range Latency
                    (micros)          (micros)          (micros)
50%                   89.00             12.00            234.00
75%                  125.00             18.00            456.00
95%                  456.00             45.00           1234.00
99%                 1847.00            124.00           4892.00
Max                45892.00           4892.00          89234.00
                       ↑ Huge max indicates hot partition!

# Check per-node load imbalance
$ nodetool status

Address      Load      Owns
10.1.0.1     547 GB    33.3%  ← Much higher than average!
10.1.0.2     238 GB    33.2%
10.1.0.3     251 GB    33.5%

# Node 10.1.0.1 has hot partition!

Fixing Hot Partitions

Solutions

1. Add Bucketing to Partition Key

# Before: Large partition
PRIMARY KEY ((celebrity_id), follower_id)

# After: Split into buckets
PRIMARY KEY ((celebrity_id, bucket), follower_id)

# Example with 10 buckets:
INSERT INTO followers (celebrity_id, bucket, follower_id, ...)
VALUES (obama_uuid, follower_id % 10, follower_id, ...)

# Now 130M followers split across 10 partitions:
# Partition 1: 13M followers
# Partition 2: 13M followers
# ... (Much more manageable!)

2. Time-Based Bucketing

# Add date to partition key
PRIMARY KEY ((user_id, date), timestamp)

# Each day = new partition
# Prevents partition from growing forever
# Old data can be archived/deleted

3. Redesign Schema (Best Solution)

# Reverse the relationship
# Query by follower instead of celebrity
PRIMARY KEY ((follower_id), celebrity_id)

# Or create separate table for hot entities
CREATE TABLE celebrity_followers_summary (
    celebrity_id uuid PRIMARY KEY,
    follower_count counter,
    top_100_followers list
);

✅ Partition Design Best Practices

1️⃣

Keep Partitions Small

Target: < 100MB per partition
Max: < 100,000 rows
Why: Fits in memory, fast queries, stable performance

2️⃣

Choose High Cardinality Keys

Good: user_id (millions of values)
Bad: status (only 3 values)
Why: More unique keys = better distribution

3️⃣

Match Query Patterns

Design: Partition key = what you query by
Example: Query by user? Use user_id
Why: Single-partition queries are fastest

4️⃣

Use Time Bucketing

Add: date or hour to partition key
Example: (user_id, date)
Why: Prevents unbounded growth

5️⃣

Avoid Composite Partitions

Use: Only when necessary
Tradeoff: More partitions = more complexity
When: Splitting hot partitions only

6️⃣

Monitor Partition Sizes

Check: nodetool tablehistograms weekly
Alert: Partitions > 100MB
Action: Redesign schema if needed

Golden Rules

  • 1 Query = 1 Partition: Design for single-partition queries
  • High Cardinality: Millions of unique partition keys
  • Bounded Growth: Use time-bucketing for time-series data
  • Test at Scale: Load test with realistic data volumes
  • Monitor Distribution: Watch for skewed load across nodes

💼 Top 10 Interview Questions

1
What is a partition key and why is it important?
+

Answer:

The partition key is the first part of the PRIMARY KEY that determines which node stores the data.

How it works:

  • Cassandra hashes the partition key value using Murmur3
  • Hash produces a token (-2^63 to 2^63-1)
  • Token determines which node owns the data
  • All rows with same partition key stored together on same node

Why it's important:

  • Performance: Single-partition queries are fastest (one node lookup)
  • Distribution: Determines how data spreads across cluster
  • Scalability: Good design enables horizontal scaling
  • Co-location: Related data stored together for efficiency

Example:

CREATE TABLE users (
    user_id uuid PRIMARY KEY,  -- partition key
    name text,
    email text
);

-- Query: SELECT * FROM users WHERE user_id = ?
-- Fast! Single node lookup
2
What's the difference between partition key and clustering key?
+

Answer:

Partition Key:

  • Determines which node stores the data
  • First part of PRIMARY KEY (in parentheses)
  • All rows with same partition key on same node
  • Used for data distribution

Clustering Key:

  • Determines sort order within a partition
  • Second part of PRIMARY KEY (after partition key)
  • Allows multiple rows per partition
  • Used for data ordering

Example:

CREATE TABLE user_events (
    user_id uuid,        -- partition key
    event_time timestamp, -- clustering key
    event_type text,
    data text,
    PRIMARY KEY ((user_id), event_time)
) WITH CLUSTERING ORDER BY (event_time DESC);

-- Partition key: user_id
--   → Determines node (all events for user on same node)
-- Clustering key: event_time DESC
--   → Sorts events newest first within partition

-- Query: SELECT * FROM user_events WHERE user_id = ? LIMIT 10
-- Returns: 10 most recent events for user (fast!)

Visual:

Partition: user_id = alice
├── event_time: 2024-12-26 10:30:00 (newest, returned first)
├── event_time: 2024-12-26 10:25:00
├── event_time: 2024-12-26 10:20:00
└── ... (sorted by clustering key)
3
How does Cassandra ensure even data distribution?
+

Answer: Cassandra uses consistent hashing with Murmur3 hash function:

Process:

  • Step 1: Hash partition key with Murmur3
    • Input: partition key value (e.g., user_id=12345)
    • Output: 64-bit token (e.g., -8234719023847192847)
  • Step 2: Map token to node
    • Each node owns range of tokens
    • Token falls in one node's range
  • Step 3: Store data on that node

Why it works:

  • Uniform Distribution: Murmur3 spreads values evenly across token range
  • Deterministic: Same key always produces same token
  • Avalanche Effect: Small key change → completely different token
  • Equal Ranges: Each node owns equal portion of token ring

Example (3-node cluster):

Token Range: -9,223,372,036,854,775,808 to 9,223,372,036,854,775,807

Node 1 owns: -9,223,372,036,854,775,808 to -3,074,457,345,618,258,603
Node 2 owns: -3,074,457,345,618,258,603 to  3,074,457,345,618,258,602  
Node 3 owns:  3,074,457,345,618,258,602 to  9,223,372,036,854,775,807

Each node owns ~33.3% of token range
→ Each stores ~33.3% of data (even distribution!)
4
What is a hot partition and how do you fix it?
+

Answer:

Hot Partition = partition that's too large or too frequently accessed

Signs of hot partition:

  • Partition size > 100MB (or > 100,000 rows)
  • Slow queries (> 100ms)
  • High read/write traffic to one node
  • Node experiencing OOM errors or GC pauses
  • Uneven load distribution across cluster

Common Causes:

  • Poor Partition Key: Low cardinality (few unique values)
  • Unbounded Growth: No time-based bucketing for time-series
  • Celebrity Problem: One entity with millions of relationships

Solutions:

1. Add Bucketing:

-- Before (hot partition)
PRIMARY KEY ((celebrity_id), follower_id)
-- 130M followers in ONE partition!

-- After (distributed across buckets)
PRIMARY KEY ((celebrity_id, bucket), follower_id)
-- Split into 10 partitions of 13M each

2. Time-Based Bucketing:

-- Add date to partition key
PRIMARY KEY ((user_id, date), timestamp)
-- Each day = new partition (bounded growth)

3. Redesign Schema:

-- Query by follower instead of celebrity
PRIMARY KEY ((follower_id), celebrity_id)
-- Now each user has small partition
5
How do you choose a good partition key?
+

Answer: Follow these criteria for partition key selection:

1. High Cardinality:

  • Many unique values (millions, not thousands)
  • ✅ Good: user_id, order_id, session_id
  • ❌ Bad: status (active/inactive), country (195 values)

2. Even Distribution:

  • Values spread evenly across all possible keys
  • ✅ Good: UUID, user_id (millions of users)
  • ❌ Bad: timestamp_hour (only 24 values)

3. Match Query Pattern:

  • Partition key should be what you query by
  • ✅ Good: Query by user? Use user_id
  • ❌ Bad: Query by user but partition by date

4. Bounded Size:

  • Partition won't grow indefinitely
  • ✅ Good: (user_id, date) for time-series
  • ❌ Bad: (user_id) for append-only data

Decision Tree:

1. What do I query by?
   → Use that as partition key

2. Will partitions grow too large?
   YES → Add time-bucketing (date, hour)
   NO  → Keep simple partition key

3. Do I have hot entities (celebrities)?
   YES → Add bucketing or redesign
   NO  → Current design OK

4. How many unique values?
   < 1,000   → Too low, rethink design
   1,000 - 1M → OK
   > 1M      → Excellent!
6
What's the difference between simple and composite partition keys?
+

Answer:

Simple Partition Key:

  • Single column determines node
  • Syntax: PRIMARY KEY (user_id)
  • Use when: One natural identifier exists
CREATE TABLE users (
    user_id uuid PRIMARY KEY,
    name text,
    email text
);
-- Partition key: user_id only

Composite Partition Key:

  • Multiple columns together determine node
  • Syntax: PRIMARY KEY ((col1, col2))
  • Use when: Need to split large partitions
CREATE TABLE sensor_data (
    sensor_id uuid,
    date text,         -- '2024-12-26'
    hour int,
    temperature decimal,
    PRIMARY KEY ((sensor_id, date), hour)
);
-- Partition key: sensor_id + date (composite)
-- One partition per sensor per day

Key Differences:

Aspect Simple Composite
Columns 1 2+
Use Case Natural key Split partitions
Query WHERE user_id=? WHERE sensor_id=? AND date=?
Complexity Simple More complex

When to use composite:

  • Time-series data (add date to avoid unbounded growth)
  • Hot partitions (add bucketing column)
  • Natural multi-column identifier
7
How does Cassandra determine which node stores a partition?
+

Answer: 4-step process using consistent hashing:

Step 1: Hash the Partition Key

Partition key: user_id = 12345
↓ Apply Murmur3 hash function
Token: -8234719023847192847 (64-bit integer)

Step 2: Locate Token in Ring

Token range: -2^63 to 2^63-1
Cluster has 3 nodes:

Node 1: -9223372036854775808 to -3074457345618258603
Node 2: -3074457345618258603 to  3074457345618258602
Node 3:  3074457345618258602 to  9223372036854775807

Token -8234719023847192847 falls in Node 1's range!

Step 3: Determine Replicas

With RF=3:
Primary replica: Node 1 (owns token range)
Replica 2: Node 2 (next node clockwise on ring)
Replica 3: Node 3 (next node clockwise on ring)

All 3 nodes store copy of this partition!

Step 4: Client Contacts Coordinator

Client driver:
1. Calculates token locally
2. Knows which nodes own that token
3. Sends request directly to one of the replica nodes
4. That node becomes coordinator for this request

Example Query Flow:

Client: SELECT * FROM users WHERE user_id = 12345

1. Driver hashes 12345 → token -8234719023847192847
2. Driver knows Node 1, 2, 3 own this token
3. Driver picks Node 1 (token-aware routing)
4. Node 1 reads data locally (same node!)
5. Returns result to client

Latency: <1ms (no network hops between nodes!)
8
What's the maximum recommended partition size?
+

Answer:

Recommended Limits:

  • Size: < 100 MB per partition
  • Rows: < 100,000 rows per partition
  • Hard Limit: 2 GB (but should never reach this!)

Why these limits?

  • Memory: Large partitions may not fit in memory
  • Compaction: Large partitions slow down compaction
  • Queries: Reading entire partition becomes slow
  • Repair: Large partitions take longer to repair
  • GC Pressure: Large objects cause GC pauses

Performance Impact:

Size Query Time Status
< 10 MB < 1 ms ✅ Excellent
10-100 MB 1-10 ms ✅ Good
100-500 MB 10-100 ms ⚠️ Warning
> 500 MB > 100 ms ❌ Bad

How to check partition sizes:

$ nodetool tablehistograms mykeyspace.users

Partition Size:
50%: 4096 bytes      (4 KB - excellent!)
75%: 16384 bytes     (16 KB - excellent!)
95%: 65536 bytes     (64 KB - good)
99%: 524288 bytes    (512 KB - acceptable)
Max: 104857600 bytes (100 MB - at limit!)
                     ↑ Watch this!
9
Can you query without specifying the partition key?
+

Answer: Technically yes, but NOT RECOMMENDED in production!

Option 1: ALLOW FILTERING (Slow!)

-- Query by non-partition-key column
SELECT * FROM users WHERE email = 'alice@example.com' 
ALLOW FILTERING;

❌ Problem:
- Scans ALL partitions on ALL nodes
- Reads entire table (could be billions of rows!)
- Can timeout or crash cluster
- Use ONLY for small tables or development

Option 2: Secondary Index (Limited Use)

CREATE INDEX ON users(email);

SELECT * FROM users WHERE email = 'alice@example.com';

⚠️ Caveats:
- Still queries ALL nodes
- Only efficient for low-cardinality columns
- Not recommended for high-cardinality (like email)
- Performance degrades with cluster size

Option 3: Materialized View (Better!)

-- Create view with different partition key
CREATE MATERIALIZED VIEW users_by_email AS
    SELECT * FROM users
    WHERE email IS NOT NULL AND user_id IS NOT NULL
    PRIMARY KEY (email, user_id);

-- Now can query efficiently by email!
SELECT * FROM users_by_email 
WHERE email = 'alice@example.com';

✅ Benefits:
- Automatically maintained by Cassandra
- Fast queries (single partition)
- Proper data distribution

Best Practice:

  • Design schema for your queries - partition key should match query pattern
  • Use materialized views for alternative query patterns
  • Never use ALLOW FILTERING in production on large tables
  • Denormalize data if you need multiple query patterns
10
How do you design partitioning for time-series data?
+

Answer: Use time-bucketing to prevent unbounded partition growth!

Problem with Naive Design:

-- BAD: Unbounded growth
CREATE TABLE sensor_readings (
    sensor_id uuid,
    timestamp timestamp,
    temperature decimal,
    humidity decimal,
    PRIMARY KEY ((sensor_id), timestamp)
);

❌ Problems:
- Partition grows forever (1 reading/second = 31M rows/year!)
- Eventually exceeds 100MB limit
- Queries slow down over time
- Compaction becomes difficult

Solution: Add Time Bucket to Partition Key

Option 1: Daily Buckets (Most Common)

CREATE TABLE sensor_readings (
    sensor_id uuid,
    date text,            -- '2024-12-26'
    timestamp timestamp,
    temperature decimal,
    humidity decimal,
    PRIMARY KEY ((sensor_id, date), timestamp)
) WITH CLUSTERING ORDER BY (timestamp DESC);

✅ Benefits:
- One partition per sensor per day
- Max 86,400 rows (1/second for 24 hours)
- Partition size: < 10 MB typically
- Old data can be TTL'd or archived

-- Query last 24 hours:
SELECT * FROM sensor_readings
WHERE sensor_id = ? AND date = '2024-12-26';

Option 2: Hourly Buckets (High Frequency)

CREATE TABLE sensor_readings (
    sensor_id uuid,
    date text,
    hour int,             -- 0-23
    timestamp timestamp,
    temperature decimal,
    PRIMARY KEY ((sensor_id, date, hour), timestamp)
);

✅ Use when:
- Very high frequency (>1/second)
- Need very small partitions
- Real-time queries on recent data

Max rows per partition: 3,600 (1/second for 1 hour)

Option 3: Weekly Buckets (Low Frequency)

CREATE TABLE sensor_readings (
    sensor_id uuid,
    week text,            -- '2024-W52'
    timestamp timestamp,
    temperature decimal,
    PRIMARY KEY ((sensor_id, week), timestamp)
);

✅ Use when:
- Low frequency data
- Long-term storage
- Historical analysis

Max rows per partition: 604,800 (1/second for 1 week)

Choosing Bucket Size:

Frequency Bucket Max Rows
> 10/sec Hourly 36,000
1-10/sec Daily 864,000
< 1/sec Weekly 604,800

Best Practices:

  • Target < 100,000 rows per partition
  • Add TTL for automatic cleanup
  • Use CLUSTERING ORDER BY timestamp DESC for recent data first
  • Query with date range for efficiency
Advertisement

Responsive Ad