Cluster Design
Architecture Design
Design scalable, resilient Cassandra deployments!
Scenario #1
High Complexity
๐ Multi-Datacenter Deployment for Global App
Your e-commerce platform is expanding globally. You have users in North America, Europe, and Asia. Design a multi-datacenter Cassandra deployment that provides low latency reads for all regions while ensuring data consistency and disaster recovery.
Requirements
- Users in 3 regions: US-East, EU-West, Asia-Pacific
- Each region needs <50ms read latency
- Writes must be durable (survive datacenter failure)
- 1M writes/day, 10M reads/day
- 99.99% availability target
โ ๏ธ Constraints
- Budget: Can deploy 6-18 nodes total
- Network latency: USโEU: 80ms, USโAsia: 150ms, EUโAsia: 200ms
- Data must be replicated to at least 2 datacenters
- No vendor lock-in (cloud-agnostic)
๐๏ธ Multi-DC Architecture Design
๐ Recommended Topology
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โ US-East DC โ โ EU-West DC โ โ Asia-Pac DC โ
โ (datacenter1) โ โ (datacenter2) โ โ (datacenter3) โ
โโโโโโโโโโโโโโโโโโโค โโโโโโโโโโโโโโโโโโโค โโโโโโโโโโโโโโโโโโโค
โ 6 nodes โโโโโบโ 6 nodes โโโโโบโ 6 nodes โ
โ RF = 3 โ โ RF = 3 โ โ RF = 3 โ
โ vnodes: 256 โ โ vnodes: 256 โ โ vnodes: 256 โ
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โฒ โฒ โฒ
โ โ โ
LOCAL_QUORUM LOCAL_QUORUM LOCAL_QUORUM
(2/3 local) (2/3 local) (2/3 local)
Total: 18 nodes (6 per DC), RF=3 per DC = 9 copies of data total
๐ง Keyspace Configuration
CREATE KEYSPACE ecommerce
WITH REPLICATION = {
'class': 'NetworkTopologyStrategy',
'datacenter1': 3, -- US-East: 3 replicas
'datacenter2': 3, -- EU-West: 3 replicas
'datacenter3': 3 -- Asia-Pac: 3 replicas
};
-- Each DC has full copy of data (3 replicas per DC)
-- Can lose any entire DC and still operate
๐ป Application Configuration
// Reads: LOCAL_QUORUM (fast, local)
consistency_level_read = LOCAL_QUORUM;
// Requires 2/3 nodes in LOCAL datacenter
// Latency: <50ms (all local)
// Writes: LOCAL_QUORUM (durable + fast)
consistency_level_write = LOCAL_QUORUM;
// Async replication to other DCs
// Latency: <50ms for client
// Connection per datacenter
local_dc = 'datacenter1'; // Set based on app location
remote_dcs = ['datacenter2', 'datacenter3'];
โก Latency Calculation
| Operation | Region | Latency | Explanation |
|---|---|---|---|
| Read | US-East | 10-15ms | LOCAL_QUORUM in US-East DC |
| Read | EU-West | 10-15ms | LOCAL_QUORUM in EU-West DC |
| Read | Asia-Pac | 10-15ms | LOCAL_QUORUM in Asia-Pac DC |
| Write | Any | 15-25ms | LOCAL_QUORUM + async to other DCs |
๐ฏ Trade-offs Analysis
| Approach | Pros | Cons |
|---|---|---|
| Chosen: RF=3 per DC | โข Full redundancy per DC โข Fast local reads โข Survive DC failure |
โข Higher storage cost โข 9 copies of data |
| Alt: RF=2 per DC | โข Lower storage cost โข Faster writes |
โข Can't lose 1 node per DC โข Less resilient |
| Alt: Active-Passive | โข Simpler setup โข Lower cost |
โข High latency for remote users โข Manual failover |
โ
Best Practices
- Use NetworkTopologyStrategy for multi-DC (never SimpleStrategy)
- RF=3 per datacenter for production (balance cost/availability)
- LOCAL_QUORUM for reads and writes (fast + consistent locally)
- Deploy app servers in same DC as Cassandra nodes
- Monitor cross-DC replication lag (should be <5 seconds)
- Test DC failover scenarios regularly
- Use snitch that matches your topology (EC2MultiRegionSnitch for AWS)
Scenario #2
Medium Complexity
๐ Cluster Sizing for IoT Workload
You're building an IoT platform that ingests sensor data. Estimate cluster size needed and design the architecture to handle the expected load with room for growth.
Workload Specs
- 10,000 sensors sending data every 30 seconds
- Each reading = 500 bytes (sensor_id, timestamp, 5 metrics)
- Data retention: 30 days raw data, 1 year aggregates
- Replication Factor: 3
- Expected growth: 2x per year
โ ๏ธ Constraints
- SLA: P99 write latency <20ms, P99 read latency <50ms
- Availability: 99.9% (can tolerate 1 node down)
- Budget: Optimize for cost while meeting SLAs
- Nodes: 16 cores, 128GB RAM, 2TB NVMe SSD each
๐ Capacity Planning & Architecture
๐ Step 1: Calculate Data Volume
Write Rate:
10,000 sensors ร (1 write / 30 sec) = 333 writes/sec
Daily Data:
333 writes/sec ร 86,400 sec/day = 28.8M writes/day
28.8M ร 500 bytes = 14.4 GB/day raw data
30-Day Storage (Raw):
14.4 GB/day ร 30 days = 432 GB raw
432 GB ร RF 3 = 1.3 TB with replication
Yearly Storage (with aggregates):
432 GB (30d raw) + 200 GB (1y aggregates) = 632 GB
632 GB ร RF 3 = 1.9 TB with replication
2-Year Projection (2x growth/year):
Year 1: 1.9 TB
Year 2: 3.8 TB (2x growth)
๐ฅ๏ธ Step 2: Determine Node Count
// Storage capacity per node (conservative)
// - Disk: 2TB
// - Usable: 50% (compaction, repairs, overhead) = 1TB
// - Safe limit: 70% = 700GB per node
Minimum nodes (storage):
Year 1: 1.9 TB / 700 GB = 3 nodes
Year 2: 3.8 TB / 700 GB = 6 nodes
Minimum nodes (availability with RF=3):
3 nodes minimum for RF=3
Write throughput check:
// Conservative: 5,000 writes/sec per node
333 writes/sec / 5,000 = 0.07 nodes
// Not a bottleneck
Recommended: Start with 6 nodes
// - Handles Year 2 projection
// - Room for spikes and maintenance
// - Better distribution with vnodes
๐๏ธ Recommended Architecture
Initial Deployment (6 Nodes):
Node1 Node2 Node3 Node4 Node5 Node6
โผ โผ โผ โผ โผ โผ
โโโโโ โโโโโ โโโโโ โโโโโ โโโโโ โโโโโ
โ C โ โ C โ โ C โ โ C โ โ C โ โ C โ
โโโโโ โโโโโ โโโโโ โโโโโ โโโโโ โโโโโ
โ โ โ โ โ โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Single Datacenter (RF=3) โ
โ Each node: 16 cores, 128GB, 2TB โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
Configuration:
- Replication Factor: 3
- vnodes: 256 per node
- Consistency: LOCAL_QUORUM (reads and writes)
- Compaction: TimeWindowCompactionStrategy (TWCS)
- TTL: 30 days on raw data table
๐ฐ Cost-Performance Analysis
| Config | Nodes | Cost/Mo | Capacity | Recommendation |
|---|---|---|---|---|
| Minimal | 3 | $900 | 1.3 TB | Too small, no growth room |
| Recommended | 6 | $1,800 | 3.8 TB | Best balance - handles 2yr growth |
| Over-provisioned | 12 | $3,600 | 7.6 TB | Unnecessary for current needs |
๐ Growth Strategy
// Scaling timeline
Year 1 (Months 1-12):
- Start: 6 nodes (3.8 TB capacity)
- Usage: ~1.9 TB (50% capacity)
- Status: โ
Good headroom
Year 2 (Months 13-24):
- Expected: 3.8 TB data (2x growth)
- At month 18: Add 3 nodes โ 9 nodes total
- New capacity: 5.7 TB
- Status: โ
Still ahead of growth
Monitoring Triggers:
- 60% disk usage โ Start planning expansion
- 70% disk usage โ Add nodes immediately
- P99 write latency > 30ms โ Add nodes for throughput
โ
Capacity Planning Best Practices
- Plan for 50% disk utilization (compaction + overhead)
- Keep 30-40% headroom for growth and spikes
- Monitor disk usage, never exceed 70-80%
- Use TWCS for time-series data (efficient compaction)
- Set TTL to auto-expire old data
- Add nodes before hitting capacity (avoid rush)
- Scale horizontally (add nodes) vs vertically (bigger nodes)
Scenario #3
Medium Complexity
๐ Choosing Replication Strategy
You're migrating from a single-DC deployment to multi-DC. Decide on the right replication strategy and how to transition without downtime.
Current Setup
- Single datacenter: 12 nodes
- Keyspaces use SimpleStrategy with RF=3
- 100TB of data across all tables
- 24/7 production traffic (no maintenance windows)
- Need to add: Second DC for DR + geographic distribution
๐ Replication Strategy Migration
๐ฏ Why NetworkTopologyStrategy?
| Strategy | Use Case | Multi-DC? | Recommendation |
|---|---|---|---|
| SimpleStrategy | Single DC only | โ No | Never use in production |
| NetworkTopologyStrategy | Production (single or multi-DC) | โ Yes | Always use this |
Key Point: Even for single-DC, use NetworkTopologyStrategy. It allows future expansion without complex migrations.
๐ Step-by-Step Migration (Zero Downtime)
-- Step 1: Update keyspace replication (Current DC)
ALTER KEYSPACE my_keyspace
WITH REPLICATION = {
'class': 'NetworkTopologyStrategy',
'dc1': 3 -- Same RF as before (SimpleStrategy RF=3)
};
-- Step 2: Run repair on all nodes (sync data)
nodetool repair -full
-- Step 3: Add second datacenter nodes
-- In cassandra.yaml for new nodes:
endpoint_snitch: GossipingPropertyFileSnitch
-- In cassandra-rackdc.properties:
dc=dc2
rack=rack1
-- Step 4: Bootstrap new DC nodes one by one
-- (Cassandra auto-streams data to new nodes)
-- Step 5: After all dc2 nodes joined, update keyspace
ALTER KEYSPACE my_keyspace
WITH REPLICATION = {
'class': 'NetworkTopologyStrategy',
'dc1': 3,
'dc2': 3 -- Now replicated to dc2!
};
-- Step 6: Run repair again to sync dc2 data
nodetool repair -full
-- Step 7: Update application config
// Point apps in dc2 to local_dc='dc2'
// Use LOCAL_QUORUM for both DCs
โฑ๏ธ Migration Timeline (100TB cluster)
| Phase | Duration | Risk | Notes |
|---|---|---|---|
| Change to NTS | 5 minutes | Low | Just metadata change |
| Repair DC1 | 24-48 hours | Low | Runs in background |
| Add DC2 nodes | 48-72 hours | Medium | Streams 100TB of data |
| Update keyspace RF | 5 minutes | Low | Metadata change |
| Repair DC2 | 24-48 hours | Low | Final sync |
| Total | 5-7 days | Low | Zero downtime! |
๐ฏ Post-Migration Configuration
Before (SimpleStrategy RF=3):
โโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Single Datacenter โ
โ 12 nodes โ
โ RF = 3 โ
โ SimpleStrategy โ โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโ
After (NetworkTopologyStrategy):
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โ DC1 (Primary) โโโโโบโ DC2 (DR/Read) โ
โ 12 nodes โ โ 12 nodes โ
โ RF = 3 โ โ RF = 3 โ
โ NTS โ
โ โ NTS โ
โ
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โ
Replication Best Practices
- ALWAYS use NetworkTopologyStrategy (even single DC)
- Never use SimpleStrategy in production
- RF=3 per datacenter is standard for production
- Run repair after any replication changes
- Throttle streaming during datacenter expansion (stream_throughput)
- Test failover between datacenters regularly
- Monitor replication lag with nodetool netstats
Scenario #4
Medium Complexity
๐ข Rack Awareness for High Availability
Your datacenter has 3 racks. Design the Cassandra deployment to survive rack failures and avoid placing replicas in the same rack.
Datacenter Setup
- 1 datacenter with 3 physical racks
- Each rack has independent power, network
- 9 nodes total: 3 nodes per rack
- RF=3 (want 1 replica per rack)
- Must survive loss of any single rack
๐ข Rack-Aware Architecture
๐๏ธ Physical Topology
Datacenter: DC1
โโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโ
โ Rack A โ โ Rack B โ โ Rack C โ
โโโโโโโโโโโโโโโโค โโโโโโโโโโโโโโโโค โโโโโโโโโโโโโโโโค
โ Node1 โ โ โ Node4 โ โ โ Node7 โ โ
โ Node2 โ โ โ Node5 โ โ โ Node8 โ โ
โ Node3 โ โ โ Node6 โ โ โ Node9 โ โ
โโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโ
โก โโโ โก โโโ โก โโโ
Independent Independent Independent
Power/Net Power/Net Power/Net
With RF=3 and 3 racks:
โ Each replica goes to different rack โ
โ Survive any single rack failure โ
๐ง Configuration (GossipingPropertyFileSnitch)
# cassandra.yaml (ALL nodes)
endpoint_snitch: GossipingPropertyFileSnitch
# cassandra-rackdc.properties (Rack A nodes - Node1,2,3)
dc=dc1
rack=rack_a
# cassandra-rackdc.properties (Rack B nodes - Node4,5,6)
dc=dc1
rack=rack_b
# cassandra-rackdc.properties (Rack C nodes - Node7,8,9)
dc=dc1
rack=rack_c
# Verify with nodetool
nodetool status
# Should show:
# DC: dc1 Rack: rack_a - Node1, Node2, Node3
# DC: dc1 Rack: rack_b - Node4, Node5, Node6
# DC: dc1 Rack: rack_c - Node7, Node8, Node9
๐ Replica Distribution Example
Data with Partition Key X (RF=3):
Primary Replica: Node2 (Rack A) โโโ
Replica 2: Node5 (Rack B) โโโผโ All different racks โ
Replica 3: Node8 (Rack C) โโโ
If Rack B fails:
โ
Still have 2/3 replicas (Rack A + Rack C)
โ
QUORUM reads/writes still work (2/3)
โ
No data loss
โ ๏ธ Common Rack Mistakes
| Configuration | Issue | Impact |
|---|---|---|
| RF=3 with 1 rack | All replicas in same rack | Rack failure = total data loss |
| RF=3 with 2 racks | 2 replicas in one rack | That rack fails = lose QUORUM |
| RF=3 with 3+ racks | 1 replica per rack | Can lose any rack โ |
โ
Rack Awareness Best Practices
- Number of racks โฅ Replication Factor (RF=3 needs 3+ racks)
- Use GossipingPropertyFileSnitch for rack awareness
- Configure cassandra-rackdc.properties correctly on each node
- Verify rack distribution with nodetool status
- Balance nodes evenly across racks (3 racks = 3/6/9/12 nodes)
- Test rack failure scenarios (nodetool disablebinary on all nodes in rack)
- Monitor rack distribution with OpsCenter or similar tools
Scenario #5
High Complexity
๐ฅ Disaster Recovery Strategy
Design a disaster recovery plan for a mission-critical application. Must handle datacenter loss, support point-in-time recovery, and meet RPO/RTO requirements.
DR Requirements
- RPO (Recovery Point Objective): <5 minutes of data loss
- RTO (Recovery Time Objective): <30 minutes to restore service
- Must survive complete loss of primary datacenter
- Need point-in-time recovery for accidental deletes/corruptions
- Compliance: 7-day backup retention
๐ฅ Disaster Recovery Architecture
๐๏ธ Multi-Layer DR Strategy
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ DR Architecture โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
Layer 1: Active-Active Multi-DC (Datacenter Failure)
โโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโ
โ Primary โโโโโโโโโโบโ Secondary โ
โ DC (US) โ Sync โ DC (EU) โ
โ RF = 3 โ <5min โ RF = 3 โ
โโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโ
RPO: <5min โ
RTO: <30min โ
Layer 2: Snapshots (Point-in-Time Recovery)
โโโโโโโโโโโโโโโโ
โ Snapshots โ โ nodetool snapshot (daily)
โ S3/Blob โ Retention: 7 days
โโโโโโโโโโโโโโโโ RPO: 24 hours, RTO: 2-4 hours
Layer 3: Backups (Long-term Retention)
โโโโโโโโโโโโโโโโ
โ Backups โ โ Incremental backups
โ S3 Glacier โ Retention: 90 days
โโโโโโโโโโโโโโโโ For compliance/audit
๐ Layer 1: Active-Active Configuration
-- Multi-DC setup
CREATE KEYSPACE production
WITH REPLICATION = {
'class': 'NetworkTopologyStrategy',
'us_east': 3, -- Primary DC
'eu_west': 3 -- DR DC (also serves EU traffic)
};
-- Application config
// Writes: LOCAL_QUORUM (async replication to other DC)
// Reads: LOCAL_QUORUM (read from local DC)
// Replication lag typically: 1-5 seconds
-- Failover procedure (if us_east fails):
// 1. DNS failover to eu_west (automated, 2-5 min)
// 2. eu_west handles all traffic with LOCAL_QUORUM
// 3. RTO: <30 minutes โ
// 4. RPO: <5 minutes (replication lag) โ
๐ธ Layer 2: Snapshot Strategy
# Daily snapshot script (cron)
#!/bin/bash
# Take snapshot (instant, uses hardlinks)
nodetool snapshot --tag daily-$(date +%Y%m%d)
# Upload to S3
for snapshot in /var/lib/cassandra/data/*/snapshots/daily*
do
aws s3 sync $snapshot s3://cassandra-backups/snapshots/
done
# Clear old snapshots (keep 7 days)
nodetool clearsnapshot -t daily-$(date -d "7 days ago" +%Y%m%d)
# Restore from snapshot:
# 1. Download from S3
# 2. Copy to data directory
# 3. nodetool refresh keyspace table
# 4. RTO: 2-4 hours (depends on data size)
๐ DR Testing Schedule
| Test Type | Frequency | Duration | Validates |
|---|---|---|---|
| DC Failover | Quarterly | 2-4 hours | Active-Active works, RTO met |
| Snapshot Restore | Monthly | 4-8 hours | Backups valid, restore process |
| Full DR Drill | Annually | 1-2 days | Complete DR playbook |
| Replication Lag | Continuous | - | RPO <5min maintained |
๐ฏ Cost-Benefit Analysis
| DR Strategy | Cost | RPO | RTO | Complexity |
|---|---|---|---|---|
| Chosen: Active-Active | High (2x nodes) | <5 min | <30 min | Medium |
| Alt: Active-Passive | Medium (1.5x nodes) | 5-15 min | 1-2 hours | Medium |
| Alt: Backups Only | Low (storage only) | 24 hours | 4-24 hours | Low |
โ
Disaster Recovery Best Practices
- Active-Active multi-DC for critical apps (best RPO/RTO)
- Take daily snapshots for point-in-time recovery
- Store backups in different region/cloud (blast radius)
- Automate failover with health checks + DNS
- Test DR procedures regularly (quarterly minimum)
- Document runbooks with exact commands
- Monitor replication lag continuously (alert if >10s)
- Use incremental backups for large clusters (faster)
Advertisement
Responsive Ad