Performance Tuning

Benchmarking in Cassandra

Measure what really matters! Benchmarking helps you test Cassandraโ€™s performance under different workloads, so you can compare results, tune the system, and ensure it delivers consistent speed and scalability in real-world conditions.

๐Ÿ“– The Story: Marcus's Misleading Benchmark

Marcus ran a benchmark that showed 500K writes/sec. He deployed to production expecting amazing performance. Instead, he got 5K writes/sec - 100x slower! His boss asked: "Why did the benchmark lie?" Marcus learned the hard way: bad benchmarks are worse than no benchmarks.

๐Ÿ˜ฑ The Bad Benchmark

What Marcus Did:

-- Marcus's "benchmark": cassandra-stress write n=10000000 \ -rate threads=1000 -- Results (single node!): Op rate: 500,000 writes/sec โ† Wow! Latency mean: 2ms โ† Amazing! Marcus: "We can handle 500K writes/sec!" โœ…

What Was Wrong:

  1. Single Node: Tested only 1 node, not cluster behavior
  2. No Replication: RF=1 (production needs RF=3!)
  3. All in Memory: Working set 200MB, server had 128GB RAM
  4. No Reads: Only writes (production is 80% reads)
  5. Sequential Keys: Perfect case, production has random keys
  6. Unrealistic Data: Fixed-size rows, production varies wildly
  7. No Compaction: Test ran 5 minutes (compaction kicks in after hours)

Production Reality:

-- Production workload (first day): Op rate: 5,000 writes/sec โ† 100x slower! ๐Ÿ’ฅ Latency P99: 250ms โ† Terrible! Timeouts: 15% โ† Users angry! -- Why so different? // - RF=3: Each write goes to 3 nodes (3x load) // - 80% reads: Mixed workload, not write-only // - Large data: 10TB doesn't fit in RAM // - Compaction: Fighting for disk I/O // - Random keys: No cache locality

The Consequences:

  • ๐Ÿ’ฅ Undersized Cluster: Planned for 2 nodes, needed 20
  • ๐Ÿ’ธ Cost Overrun: 10x more infrastructure needed
  • โฐ Timeline Blown: 3-month delay to add capacity
  • ๐Ÿ˜ก Users Unhappy: Slow, unreliable for weeks
  • ๐ŸŽฏ Boss Unhappy: "Why didn't you test properly?"

โœ… The Proper Benchmark

What Marcus Should Have Done:

-- Realistic benchmark (3-node cluster, RF=3): cassandra-stress mixed \ ratio\(write=2,read=8\) \ # 80% reads, 20% writes n=100000000 \ # 100M rows (realistic size) -pop dist=GAUSSIAN\(1..100000000,50000000,10000000\) \ -col size=FIXED\(1024\) \ # 1KB rows -rate threads=100 \ # Realistic concurrency -node node1,node2,node3 \ # All nodes -schema "replication(factor=3)" -- Run for 24+ hours to include compaction!

Realistic Results:

  • ๐Ÿ“Š Throughput: 25K ops/sec total (not 500K!)
  • โฑ๏ธ Write Latency: P99 = 15ms (not 2ms!)
  • โฑ๏ธ Read Latency: P99 = 8ms
  • ๐Ÿ’พ Disk I/O: 60% utilized (not 10%!)
  • ๐Ÿ”„ Compaction: Running continuously
  • ๐ŸŽฏ Realistic: Matches production!

Proper Capacity Planning:

  • โœ… Expected Load: 50K ops/sec
  • โœ… Benchmark Shows: 25K ops/sec per 3-node cluster
  • โœ… Calculation: Need 6 nodes (2ร— for 2x headroom)
  • โœ… Budget: Planned correctly from start
  • โœ… Result: Production launch smooth! ๐ŸŽ‰

Marcus learned: Benchmark production workload, not toy scenarios! ๐ŸŽฏ

๐ŸŽฏ Benchmarking Fundamentals

The principles of meaningful performance testing!

๐Ÿ“ The Golden Rules

  1. Match Production Workload: Same read/write ratio, data size, access patterns
  2. Match Production Cluster: Same RF, consistency levels, number of nodes
  3. Run Long Enough: 24+ hours to include compaction, GC cycles
  4. Realistic Data: Production-like data sizes and distributions
  5. Monitor Everything: CPU, disk, network, latency percentiles
  6. Don't Trust Peaks: Sustained performance matters, not bursts

Common Benchmarking Mistakes

โŒ

Bad Benchmark

  • Single node test
  • RF=1 (no replication)
  • Write-only workload
  • Sequential keys
  • 5-minute duration
  • All data in RAM
  • Report peak throughput

100x off from production!

โœ…

Good Benchmark

  • Full cluster (โ‰ฅ3 nodes)
  • RF=3 (production config)
  • Mixed read/write (80/20)
  • Random/Gaussian keys
  • 24+ hour duration
  • Data > RAM (realistic)
  • Report sustained P99 latency

Predicts production!

Metrics That Matter

Metric Why It Matters What to Report
Throughput How many ops/sec sustained Sustained rate, not peak
P99 Latency Worst-case user experience 99th percentile (not average!)
Resource Util Bottlenecks and headroom CPU, disk, network %
GC Pauses Impact on tail latency Max pause time, frequency
Compaction Long-term stability Pending tasks, SSTable count

Don't Trust Marketing Benchmarks!

Vendor benchmarks often show:

  • Peak throughput (not sustained)
  • Best-case scenarios (all in memory)
  • Single node (not distributed)
  • Synthetic workloads (not production-like)

Always run your own benchmarks with your workload!

๐Ÿ”จ cassandra-stress Tool

The official Cassandra benchmarking tool!

Basic Usage

Simple Write Test

-- Write 1 million rows: cassandra-stress write n=1000000 -- Specify node: cassandra-stress write n=1000000 \ -node 192.168.1.10 -- Multiple nodes (load balanced): cassandra-stress write n=1000000 \ -node 192.168.1.10,192.168.1.11,192.168.1.12

Read Test

-- Read 1 million operations: cassandra-stress read n=1000000 \ -node 192.168.1.10 -- Note: Must write data first! cassandra-stress write n=10000000 -node 192.168.1.10 cassandra-stress read n=1000000 -node 192.168.1.10

Mixed Workload (Most Important!)

-- 80% reads, 20% writes (typical production): cassandra-stress mixed \ ratio\(write=2,read=8\) \ n=10000000 \ -node 192.168.1.10,192.168.1.11,192.168.1.12 -- Different ratios: ratio\(write=1,read=9\) # 90% reads, 10% writes ratio\(write=3,read=7\) # 70% reads, 30% writes ratio\(write=5,read=5\) # 50/50 split

Advanced Options

#################################### # PRODUCTION-LIKE BENCHMARK #################################### cassandra-stress mixed \ ratio\(write=2,read=8\) \ n=100000000 \ # 100M operations \ # Data distribution (CRITICAL!): -pop dist=GAUSSIAN\(1..100000000,50000000,10000000\) \ # โ†‘ Gaussian distribution around 50M with 10M stddev # (Realistic: some keys accessed more than others) \ # Column size: -col size=FIXED\(1024\) \ # 1KB rows # OR variable size: # -col size=UNIFORM(100..10000) \ # 100 bytes to 10KB \ # Concurrency: -rate threads=100 \ # 100 concurrent threads \ # Throttling (optional): # -rate threads=100 limit=50000/s \ # Max 50K ops/sec \ # Cluster nodes: -node node1,node2,node3 \ \ # Schema (replication!): -schema "replication(factor=3)" \ \ # Consistency level: -mode native cql3 \ -log file=benchmark.log

Access Patterns

โŒ

Sequential (Unrealistic)

-pop seq=1..1000000 /* Keys: 1, 2, 3, 4... Perfect cache locality NOT like production! */
โš ๏ธ

Uniform (Somewhat Realistic)

-pop dist=UNIFORM\(1..1000000\) /* All keys equal probability Better than sequential */
โœ…

Gaussian (Realistic!)

-pop dist=GAUSSIAN\(1..1M,500K,100K\) /* Hot keys + cold keys Matches production! */

Reading Results

-- Example output: /* Results: Op rate : 45,237 op/s [WRITE: 9,047, READ: 36,190] Partition rate : 45,237 pk/s [WRITE: 9,047, READ: 36,190] Row rate : 45,237 row/s [WRITE: 9,047, READ: 36,190] Latency mean : 2.2 ms [WRITE: 2.8 ms, READ: 2.0 ms] Latency median : 1.8 ms [WRITE: 2.3 ms, READ: 1.7 ms] Latency 95th percentile : 4.1 ms [WRITE: 5.2 ms, READ: 3.8 ms] Latency 99th percentile : 8.7 ms [WRITE: 11.3 ms, READ: 7.9 ms] Latency 99.9th percentile : 28.4 ms [WRITE: 35.2 ms, READ: 25.1 ms] Latency max : 247.2 ms [WRITE: 247.2 ms, READ: 189.3 ms] Total operation time : 00:36:52 */ -- What to look at: โœ… Op rate: Sustained throughput (not peak!) โœ… P99 latency: User experience (not mean!) โœ… Max latency: Worst case observed โš ๏ธ If P99 > 50ms โ†’ performance problem!

๐Ÿ“‹ Custom Workload Profiles

Testing your actual schema and queries!

Creating Custom Schema

-- myworkload.yaml: # # Custom schema matching production # keyspace: mykeyspace keyspace_definition: | CREATE KEYSPACE mykeyspace WITH replication = { 'class': 'NetworkTopologyStrategy', 'datacenter1': 3 }; table: users table_definition: | CREATE TABLE users ( user_id uuid PRIMARY KEY, username text, email text, created_at timestamp, last_login timestamp, profile_data blob ); columnspec: - name: user_id population: uniform(1..10000000) - name: username size: uniform(5..20) - name: email size: uniform(10..50) - name: profile_data size: uniform(1024..10240) # 1-10KB insert: partitions: fixed(1) batchtype: UNLOGGED queries: read1: cql: SELECT * FROM users WHERE user_id = ? fields: samerow read2: cql: SELECT username, email FROM users WHERE user_id = ? fields: samerow

Running Custom Workload

-- Use custom profile: cassandra-stress user \ profile=myworkload.yaml \ ops\(insert=1\) \ n=10000000 \ -node node1,node2,node3 -- Mixed operations with custom queries: cassandra-stress user \ profile=myworkload.yaml \ ops\(insert=1,read1=5,read2=4\) \ n=50000000 \ -rate threads=100 \ -node node1,node2,node3

Time-Series Workload Example

-- timeseries.yaml: keyspace: metrics table_definition: | CREATE TABLE metrics ( device_id text, timestamp timestamp, temperature double, humidity double, PRIMARY KEY (device_id, timestamp) ) WITH CLUSTERING ORDER BY (timestamp DESC) AND compaction = { 'class': 'TimeWindowCompactionStrategy', 'compaction_window_size': 1, 'compaction_window_unit': 'DAYS' }; columnspec: - name: device_id population: uniform(1..10000) # 10K devices - name: timestamp cluster: fixed(1) queries: latest: cql: SELECT * FROM metrics WHERE device_id = ? LIMIT 100 fields: samerow range: cql: SELECT * FROM metrics WHERE device_id = ? AND timestamp > ? AND timestamp < ? fields: samerow

๐Ÿ“Š Monitoring During Benchmarks

What to watch while benchmarking!

Essential Metrics

System Resources

-- Monitor CPU: top -d 1 // Target: 60-80% CPU (headroom for peaks) -- Monitor disk I/O: iostat -x 1 // Watch: %util (should be < 80%) // await (should be < 10ms) -- Monitor network: iftop -i eth0 // Check: Not saturating network bandwidth -- Monitor memory: free -h // Check: Not swapping (swap = 0)

Cassandra Metrics

-- Check compaction: nodetool compactionstats /* pending tasks: 12 โ† Should be < 20 SSTable count: 87 โ† Should be < 100 */ -- GC statistics: nodetool gcstats /* Max GC time: 157ms โ† Should be < 200ms */ -- Table statistics: nodetool tablestats mykeyspace.users // Check: Space used, read/write latency

Application Metrics

-- Cassandra-stress output: // Watch in real-time: // - Op rate (throughput) // - P99 latency (user experience) // - Error rate (should be 0%!) -- If errors appear: // Check: system.log for exceptions // Common: Timeouts, unavailable

Monitoring Script

-- monitor.sh (run during benchmark): #!/bin/bash LOGFILE="benchmark_monitor.log" while true; do echo "=== $(date) ===" >> $LOGFILE # CPU echo "CPU:" >> $LOGFILE mpstat 1 1 | tail -1 >> $LOGFILE # Disk echo "Disk:" >> $LOGFILE iostat -x 1 1 | grep -E "Device|sda|nvme" >> $LOGFILE # Memory echo "Memory:" >> $LOGFILE free -h | grep -E "Mem|Swap" >> $LOGFILE # Cassandra echo "Compaction:" >> $LOGFILE nodetool compactionstats | head -5 >> $LOGFILE echo "" >> $LOGFILE sleep 60 done

โญ Benchmarking Best Practices

Lessons learned from production benchmarks!

1. Warm Up First

Don't measure cold start performance!

-- Bad: Measure immediately cassandra-stress write n=10000000 โ† Includes cold start! -- Good: Warm up first # 1. Write data cassandra-stress write n=100000000 -node node1,node2,node3 # 2. Warm up OS cache (read 10% of data) cassandra-stress read n=10000000 -node node1,node2,node3 # 3. NOW measure cassandra-stress mixed ratio\(write=2,read=8\) \ n=50000000 \ -node node1,node2,node3

2. Run Long Enough

Short tests miss compaction impact!

-- Too short (5 minutes): // - No compaction yet // - No GC pressure // - All in memory // Result: Unrealistic! -- Minimum duration: // - Write-heavy: 24 hours (see compaction) // - Read-heavy: 4-8 hours (cache warm) // - Mixed: 12-24 hours -- Use duration parameter: cassandra-stress mixed ratio\(write=2,read=8\) \ duration=24h \ โ† Run for 24 hours! -rate threads=100 \ -node node1,node2,node3

3. Test Data Larger Than RAM

Force disk I/O like production!

-- Bad: Fits in RAM Server RAM: 128GB Dataset: 10GB โ† All cached! Unrealistic! -- Good: Exceeds RAM Server RAM: 128GB Dataset: 500GB โ† Forces disk reads! Realistic! -- Calculation: // Per node: 500GB / 3 nodes = 167GB per node // OS cache: ~50GB // Hot set: 30% cached, 70% from disk โœ…

4. Test With Production RF

RF=3 means 3x more work!

-- Wrong: cassandra-stress write n=10000000 # Default: RF=1 (no replication!) # Writes: 1x load -- Correct: cassandra-stress write n=10000000 \ -schema "replication(factor=3)" # RF=3: Each write goes to 3 nodes # Writes: 3x load!

5. Test Failure Scenarios

How does it perform when things go wrong?

-- During benchmark, simulate failures: # 1. Kill one node (test RF resilience) kill -9 # Expect: Slight latency increase, no errors # 2. Network partition iptables -A INPUT -s node2 -j DROP # Expect: CL=QUORUM still works with 2/3 nodes # 3. Disk full dd if=/dev/zero of=/var/lib/cassandra/data/fill bs=1M # Expect: Node goes read-only, cluster continues

The 80/20 Rule

For capacity planning:

  • Benchmark shows: 50K ops/sec sustained
  • Plan for: 40K ops/sec (80% of benchmark)
  • Reason: Leaves headroom for:
    • Traffic spikes
    • Node failures
    • Repairs/compaction
    • GC pauses

๐Ÿ’ผ Interview Questions & Expert Answers

Master benchmarking for your interview!

1 Why is benchmarking a single node with RF=1 misleading for production planning? โ–ผ

Answer: Single-node RF=1 benchmark doesn't account for: (1) replication overhead (RF=3 means 3x write load), (2) network latency between nodes, (3) coordinator overhead, (4) consistency level coordination, (5) cluster-wide compaction impact. Can be 10-100x off from production performance.

What Single-Node RF=1 Misses:

-- Single node, RF=1 benchmark: cassandra-stress write n=10000000 -node node1 Result: 500K writes/sec โœ… Why so fast? // 1. No replication (write once, done!) // 2. No network hops (local only) // 3. No coordinator overhead // 4. No consistency coordination // 5. All data on one disk

Production Reality (RF=3, 3 nodes):

-- Same benchmark, RF=3 cluster: cassandra-stress write n=10000000 \ -node node1,node2,node3 \ -schema "replication(factor=3)" Result: 50K writes/sec โ† 10x slower! Why slower? # 1. Replication overhead (3x writes): // Client โ†’ Coordinator // Coordinator โ†’ Replica 1 (self) // Coordinator โ†’ Replica 2 (network) // Coordinator โ†’ Replica 3 (network) // Wait for CL=QUORUM (2/3 responses) // Respond to client # 2. Network latency: // Each hop adds 0.5-2ms // 2 remote writes = 1-4ms added # 3. Coordinator CPU: // Routing, coordination, consistency # 4. Distributed compaction: // All 3 nodes compacting simultaneously // Fighting for disk I/O # 5. Load distribution: // Some nodes may be hotspots

Performance Impact Breakdown:

Factor Single Node RF=1 Cluster RF=3
Writes per operation 1 (local) 3 (distributed)
Network hops 0 2 (+ 2-4ms)
Coordination overhead None CL=QUORUM logic
Disk contention Single disk 3 disks (better!)
Compaction impact Local only Cluster-wide

Key Takeaway: Always benchmark full cluster with production RF. Single-node tests can be 10-100x faster than production reality!

2 Why should benchmarks run for 24+ hours instead of 5 minutes? โ–ผ

Answer: Short benchmarks miss critical long-term effects: (1) compaction overhead kicks in after hours, (2) GC pressure builds as heap fills, (3) OS cache effectiveness stabilizes, (4) SSTable count grows, (5) disk I/O patterns change. 5-minute test shows "best case", 24-hour shows "sustained reality".

What Happens Over Time:

0-5 Minutes (Honeymoon Phase):

-- Initial state: Throughput: 100K writes/sec โ† AMAZING! Latency P99: 5ms โ† PERFECT! Why so good? // - Everything in memory (RAM cache) // - No compaction yet (few SSTables) // - No GC pressure (heap mostly empty) // - Minimal disk I/O This is NOT sustainable!

1-4 Hours (Reality Sets In):

-- Compaction starts: Throughput: 70K writes/sec โ† 30% drop Latency P99: 15ms โ† 3x worse What changed? // - Memtables flushed (10+ SSTables now) // - Compaction started (fighting for disk I/O) // - Bloom filters loaded (using RAM) // - Some GC pauses (heap 60% full)

8-12 Hours (Steady State):

-- Sustained performance: Throughput: 50K writes/sec โ† Stabilized Latency P99: 25ms โ† This is reality! Steady state characteristics: // - Compaction running continuously // - 20-50 SSTables per table // - GC every few seconds // - Disk I/O 60-80% // - OS cache stabilized THIS is production performance! โœ…

24+ Hours (Long-Term Stability):

-- Verify stability: Throughput: 48-52K writes/sec โ† Consistent โœ… Latency P99: 22-28ms โ† Stable โœ… What you're testing: // - No performance degradation over time // - Compaction keeps up with writes // - GC pauses stay reasonable // - Disk space usage predictable // - No memory leaks

Performance Over Time Graph:

Throughput (ops/sec) 100K |โ–ˆโ–ˆโ–ˆโ–ˆ โ† 5 min: 100K 80K |โ–ˆโ–ˆโ–ˆโ–ˆ 60K | โ–ˆโ–ˆโ–ˆ โ† 2 hr: 70K 40K | โ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆโ–ˆ โ† 8+ hr: 50K (steady) 20K | 0 |โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ 0 2hr 8hr 24hr Short benchmark = Misleading 100K! Long benchmark = Reality 50K!

Key Takeaway: 5-minute tests show best-case burst performance. 24-hour tests show sustainable production reality. Always test long enough to see compaction impact!

3 What's the difference between sequential and Gaussian access patterns in benchmarks? Why does it matter? โ–ผ

Answer: Sequential (keys 1,2,3...) has perfect cache locality - all data accessed in order. Gaussian (hot keys + cold keys) mimics production where some data is accessed frequently, some rarely. Sequential can be 10-50x faster due to caching, making it unrealistic for capacity planning.

Sequential Access Pattern:

-- Sequential benchmark: cassandra-stress write n=10000000 \ -pop seq=1..10000000 Access pattern: Write: key 1, 2, 3, 4, 5... Read: key 1, 2, 3, 4, 5... Why unrealistic? // 1. Perfect cache locality: // Read key 1 โ†’ cached // Read key 2 โ†’ next block, cached // Read key 3 โ†’ same block, cached // Result: 95%+ cache hit rate! // 2. Sequential disk reads: // SSD: 3500 MB/s sequential // vs 500 MB/s random // 7x faster! // 3. No data skew: // All keys equally accessed // Production: 80/20 rule! Result: 200K reads/sec, 2ms P99 โœ… But... not realistic!

Gaussian Access Pattern:

-- Gaussian benchmark: cassandra-stress write n=10000000 \ -pop dist=GAUSSIAN\(1..10000000,5000000,1000000\) # range mean stddev Access pattern: // Bell curve centered at 5M: // Keys 4M-6M: Accessed frequently (hot data) // Keys 3M-7M: Accessed sometimes (warm data) // Keys 1M-3M, 7M-10M: Rarely accessed (cold data) Why realistic? // 1. Cache behavior: // Hot keys (20%): 90% cache hit // Warm keys (30%): 50% cache hit // Cold keys (50%): 10% cache hit // Overall: 40-50% cache hit (realistic!) // 2. Random disk reads: // Jumping around SSTables // 500 MB/s random (not 3500 MB/s!) // 3. Data skew: // Matches production 80/20 rule // 20% of data gets 80% of requests Result: 25K reads/sec, 15ms P99 โŒ But... this IS production reality!

Performance Comparison:

Pattern Cache Hit Rate Read Throughput P99 Latency
Sequential 95% 200K ops/sec 2ms
Uniform 30% 40K ops/sec 12ms
Gaussian 45% 50K ops/sec 10ms
Production 40-50% 45-55K ops/sec 8-15ms

Choosing the Right Pattern:

  • โŒ Sequential: Never use (except for debugging)
  • โœ… Gaussian: Best for most workloads (hot data + cold data)
  • โœ… Uniform: Good for evenly distributed workloads (rare)
  • โญ Custom: Best if you know your actual distribution

Key Takeaway: Sequential access is 4-10x faster due to perfect caching. Always use Gaussian or production-like distributions for realistic capacity planning!

4 Your benchmark shows 100K writes/sec. How many nodes do you need for 200K writes/sec in production with 2x headroom? โ–ผ

Answer: 8 nodes. Calculation: Need 200K ร— 2 = 400K capacity. Benchmark shows 100K on 3 nodes = 33K per node. 400K / 33K = 12 nodes. But account for RF=3 coordination overhead and node failures, so ~8-10 nodes with proper tuning.

Step-by-Step Calculation:

Step 1: Understand Your Benchmark

-- Benchmark results: Cluster: 3 nodes RF: 3 Throughput: 100K writes/sec Duration: 24 hours (sustained!) -- Per-node capacity: 100K writes/sec / 3 nodes = 33K writes/sec per node BUT WAIT! This is COORDINATOR capacity, not write capacity!

Step 2: Account for RF=3

-- Each write hits 3 nodes: Client โ†’ Coordinator (Node 1) Coordinator โ†’ Replica 1 (self) Coordinator โ†’ Replica 2 (Node 2) Coordinator โ†’ Replica 3 (Node 3) -- Actual writes per node: 100K coordinator ops ร— 3 replicas = 300K physical writes 300K / 3 nodes = 100K physical writes per node Each node is actually handling 100K writes/sec!

Step 3: Calculate for Target Load

-- Target: 200K writes/sec -- With 2x headroom: 400K capacity needed -- Naive calculation: 400K / 33K per node = 12 nodes -- But this assumes linear scaling (doesn't exist!)

Step 4: Account for Scaling Factors

-- Scaling efficiency: 3 nodes: 100% baseline 6 nodes: 90% efficiency (gossip overhead) 12 nodes: 80% efficiency (more coordination) 24 nodes: 70% efficiency -- Adjusted calculation: 6 nodes: 100K ร— 2 ร— 0.90 = 180K writes/sec 9 nodes: 100K ร— 3 ร— 0.85 = 255K writes/sec 12 nodes: 100K ร— 4 ร— 0.80 = 320K writes/sec -- With 2x headroom for 200K target: Need: 400K capacity Closest: 12 nodes (320K) โŒ Still short!

Step 5: Conservative Recommendation

-- Reality check factors: 1. Node failure: Lose 1 node = 8% capacity loss 2. Repairs: Uses 10-20% capacity 3. Compaction spikes: Temporary slowdowns 4. GC pauses: Periodic latency spikes 5. Uneven load distribution: Some nodes hotter -- Recommended approach: Start with: 8-10 nodes Reasoning: // - 8 nodes: 100K ร— 2.66 ร— 0.85 = 226K // - With headroom: 113K sustained // - Close to 200K target // - Can scale to 12 if needed Monitor and add nodes if: // - CPU > 70% // - Disk I/O > 75% // - P99 latency > 50ms

Quick Formula:

nodes_needed = (target_load ร— headroom) / (bench_throughput / bench_nodes ร— efficiency) Example: = (200K ร— 2) / (100K / 3 ร— 0.85) = 400K / 28K = 14 nodes (conservative) OR start with 8-10, scale up as needed โœ…

Key Takeaway: Always include 2x headroom, account for scaling efficiency (80-90%), and monitor to add capacity before hitting limits. Start conservative, scale up as needed!

5 What metrics matter most when evaluating benchmark results? Why is P99 latency more important than average? โ–ผ

Answer: P99 latency shows worst-case user experience (1 in 100 requests). Average hides tail latency - can have 2ms average with 500ms P99. Also critical: sustained throughput (not peak), resource utilization (headroom), GC pause times, and compaction keeping up. P99 > 50ms means poor user experience.

Why Average Latency is Misleading:

-- Example distribution: 90 requests: 2ms 5 requests: 10ms 4 requests: 50ms 1 request: 500ms โ† One really slow request! -- Calculated metrics: Average: (90ร—2 + 5ร—10 + 4ร—50 + 1ร—500) / 100 = 9.3ms P50 (median): 2ms P90: 2ms P99: 500ms โ† The truth! Marketing: "Average latency: 9ms!" โœ… Reality: "1 in 100 users wait 500ms!" โŒ

Critical Metrics Ranked:

Metric Why It Matters Target
P99 Latency โญ Worst-case user experience < 50ms
Sustained Throughput โญ Can you handle load long-term? 2x target load
Resource Utilization Headroom for spikes CPU < 70%, Disk < 75%
P999 Latency Absolute worst case < 100ms
GC Max Pause Causes latency spikes < 200ms
Compaction Pending Long-term stability < 20 tasks
Average Latency Mostly marketing Ignore this!

Real-World Example:

-- Benchmark A (looks good on paper): Throughput: 150K ops/sec โ† Great! Latency mean: 3ms โ† Amazing! Latency P99: 250ms โ† TERRIBLE! Resource util: 95% โ† No headroom! Verdict: FAIL - will have timeouts in production -- Benchmark B (looks worse on paper): Throughput: 80K ops/sec โ† Lower Latency mean: 8ms โ† Higher Latency P99: 25ms โ† EXCELLENT! โœ… Resource util: 65% โ† Good headroom! โœ… Verdict: PASS - predictable, scalable, production-ready

How to Read Benchmark Output:

-- cassandra-stress output: /* Latency mean : 8.2 ms โ† Ignore Latency median : 6.1 ms โ† Ignore Latency 95th percentile : 18.4 ms โ† Good data point Latency 99th percentile : 32.7 ms โ† CRITICAL! โญ Latency 99.9th percentile : 87.3 ms โ† Important Latency max : 247.2 ms โ† Watch for outliers */ -- Decision criteria: P99 < 30ms? โœ… Excellent P99 30-50ms? โš ๏ธ Acceptable P99 > 50ms? โŒ Needs tuning P999 > 100ms? โŒ Major problem

Key Takeaway: P99 shows real user experience. Average is marketing fluff. Always demand P99 < 50ms and P999 < 100ms for production systems!

๐ŸŽ“ Chapter Summary: Benchmarking Mastery

You now understand production-grade benchmarking!

Golden Rules:

  • ๐ŸŽฏ Full Cluster: 3+ nodes with RF=3 (not single node!)
  • โฑ๏ธ Run Long: 24+ hours to see compaction
  • ๐Ÿ“Š Match Production: 80/20 read/write, Gaussian distribution
  • ๐Ÿ’พ Data > RAM: Force disk I/O
  • ๐Ÿ“ˆ Watch P99: Not average (P99 < 50ms)
  • ๐ŸŽš๏ธ 2x Headroom: For spikes and failures

Essential Command:

cassandra-stress mixed ratio\(write=2,read=8\) \ duration=24h \ -pop dist=GAUSSIAN\(1..100M,50M,10M\) \ -rate threads=100 \ -node node1,node2,node3 \ -schema "replication(factor=3)"

Red Flags:

  • โŒ 5-minute test (misses compaction!)
  • โŒ Sequential keys (unrealistic cache)
  • โŒ Write-only (production is mixed!)
  • โŒ Reporting average latency (watch P99!)

Remember Marcus: Bad benchmarks are worse than no benchmarks! ๐ŸŽฏ

Advertisement

Responsive Ad