Multi-Node Cluster Setup
High availability, fault tolerance, and massive scale - Production ready!
🌐 Why Multi-Node Clusters?
The Power of Multiple Nodes
Multi-node clusters unlock Cassandra's true potential!
✅ Benefits of Multi-Node
- 💪 High Availability: No single point of failure
- 📊 Scalability: Add nodes to handle more data/traffic
- 🔄 Data Replication: Multiple copies for safety
- ⚡ Better Performance: Distribute workload
- 🌍 Geographic Distribution: Multi-datacenter support
- 🛡️ Fault Tolerance: Survive node failures
- 🔧 Maintenance: Rolling upgrades without downtime
❌ Single-Node vs Multi-Node
Single Node:
- ❌ Node dies = entire system down
- ❌ Disk fills = can't write data
- ❌ No redundancy = data loss risk
Multi-Node (3+ nodes):
- ✅ Lose 1 node = system still works!
- ✅ Data replicated across nodes
- ✅ Can scale horizontally
🚀 Multi-node is REQUIRED for production!
🏗️ Multi-Node Architecture
Understanding how nodes work together!
3-Node Cluster Ring
Minimum: 3 Nodes
Industry standard for production
- ✅ Survive 1 node failure
- ✅ RF=3 with quorum
- ✅ Good for small-medium apps
- 💰 Cost-effective
Typical: 5-10 Nodes
Most production clusters
- ✅ Better distribution
- ✅ More fault tolerance
- ✅ Higher throughput
- 📊 Handles most workloads
Large: 100+ Nodes
Netflix, Apple, Instagram
- 🚀 Massive scale
- 🌍 Multi-datacenter
- ⚡ Petabytes of data
- 💪 Enterprise workloads
Key Concepts
- Peer-to-peer: No master/slave - all nodes are equal
- Token ranges: Each node owns part of the ring
- Replication: Data copied to multiple nodes (RF)
- Consistency level: How many replicas must respond
- Gossip protocol: Nodes share cluster state
- Snitch: Determines datacenter and rack placement
🐳 Setting Up 3-Node Cluster with Docker
Easiest way to create a multi-node cluster!
Create docker-compose.yml
Start the Cluster
Verify Cluster Status
Cluster is Ready! 🎉
Your 3-node cluster is now running!
- ✅ 3 nodes: All UP and NORMAL
- ✅ Data distributed across nodes
- ✅ Ready for replication
- ✅ Fault tolerant (survives 1 node failure)
Be Patient During Startup!
Cluster formation takes time:
- ⏱️ Node 1: Starts first (1 minute)
- ⏱️ Node 2: Discovers node 1 (2 minutes total)
- ⏱️ Node 3: Joins cluster (3-4 minutes total)
- ✅ Wait for ALL nodes to show UN before using!
💾 Data Replication
How data is copied across nodes for safety!
Understanding Replication Factor (RF)
RF = Number of copies of your data
- RF=1: No redundancy (single-node only)
- RF=2: Minimum for production (can lose 1 node)
- RF=3: Industry standard (can lose 1 node with quorum)
- RF=5: High availability (can lose 2 nodes)
Create Keyspace with RF=3
Test Replication
SimpleStrategy
Single datacenter
- Easy to configure
- Perfect for learning
- Good for single-DC production
- Replicas placed sequentially
NetworkTopologyStrategy
Multi-datacenter
- Production standard
- Specify RF per datacenter
- Geographic distribution
- Rack-aware placement
Choosing Replication Factor
Rule of thumb:
- 3-node cluster: RF=3 (every node has all data)
- 5-node cluster: RF=3 (good balance)
- 10+ node cluster: RF=3 (still standard)
Formula: Nodes that can fail = RF - (RF/2 + 1)
- RF=3: Can lose 1 node (3 - 2 = 1)
- RF=5: Can lose 2 nodes (5 - 3 = 2)
⚖️ Consistency Levels
Balance between consistency and availability!
What is Consistency Level?
How many replicas must respond for a read/write to succeed?
Example: RF=3, Consistency Level = QUORUM
- Write succeeds when 2 of 3 replicas acknowledge
- Read gets data from 2 of 3 replicas
- Ensures you always read what you wrote!
ONE - Fastest, Least Consistent
Only 1 replica must respond
- ⚡ Fastest performance
- ⚠️ May read stale data
- 💡 Use for: Logs, metrics, non-critical data
QUORUM - Balanced (RECOMMENDED)
Majority of replicas must respond
- ⚖️ Good balance
- ✅ Strong consistency
- 💡 Use for: Most production workloads
- 📊 Formula: (RF/2) + 1
ALL - Slowest, Most Consistent
ALL replicas must respond
- 🔒 Strongest consistency
- 🐌 Slowest performance
- ❌ If 1 replica down = operation fails
- 💡 Use for: Critical financial data only
LOCAL_QUORUM - Multi-DC
Quorum within local datacenter
- 🌍 For multi-datacenter clusters
- ⚡ No cross-DC latency
- 💡 Use for: Geo-distributed apps
Using Consistency Levels in CQL
Production Recommendation
For most applications:
- ✅ RF=3 (replication factor)
- ✅ QUORUM for reads and writes
- ✅ Survives 1 node failure
- ✅ Strong consistency guarantee
- ✅ Good performance
➕ Adding and Removing Nodes
Scale your cluster dynamically!
Adding a New Node
Example: Add 4th node to 3-node cluster
Run Cleanup After Adding Nodes
Removing a Node (Decommission)
Important: Decommission, Don't Just Stop!
Always use nodetool decommission before removing a node!
- ✅ Correct: decommission → stop → remove
- ❌ Wrong: Just stop/remove (data loss!)
- ⚠️ Decommission streams data to other nodes first
- ⚠️ Without decommission: lose data permanently!
📊 Monitoring Multi-Node Clusters
Keep your cluster healthy!
Check Cluster Health
Monitor Data Distribution
Monitor Performance
Check Replication
Production Monitoring Tools
For serious production clusters:
- Prometheus + Grafana: Industry standard, free
- DataStax OpsCenter: Commercial, comprehensive
- JMX Metrics: Built-in, use with JConsole
- cAdvisor: Docker container monitoring
💡 Production Best Practices
Run clusters like the pros!
✅ Minimum 3 Nodes for Production
Never run production with < 3 nodes!
- Need RF=3 for fault tolerance
- Need quorum (2 nodes) even if 1 fails
- 3 is absolute minimum, 5+ is better
✅ Use NetworkTopologyStrategy
✅ Run Regular Repairs
✅ Take Regular Snapshots
✅ Monitor Everything
Key metrics to watch:
- Node status (all UN?)
- Disk space (< 80% full?)
- Heap memory (< 75%?)
- Read/write latency
- Compaction pending
✅ Plan for Growth
Scale proactively:
- Add nodes before disk hits 70%
- Add nodes before latency degrades
- Better to have extra capacity
- Adding nodes takes time (plan ahead!)
Common Mistakes to Avoid
- ❌ Running production with < 3 nodes
- ❌ Using RF=1 in production
- ❌ Removing nodes without decommissioning
- ❌ Never running repairs
- ❌ Not monitoring disk space
- ❌ Using ALL consistency level everywhere
- ❌ Not testing backups/restores
🎉 You're Ready for Production Clusters!
Congratulations! You now understand multi-node Cassandra clusters!
🎓 What You Learned:
- 🌐 Why multi-node: High availability, scalability, fault tolerance
- 🏗️ Architecture: Peer-to-peer, token rings, gossip protocol
- 🐳 Docker setup: 3-node cluster with Docker Compose
- 💾 Replication: RF=3, SimpleStrategy vs NetworkTopologyStrategy
- ⚖️ Consistency: ONE, QUORUM, ALL - choose wisely
- ➕ Scaling: Add/remove nodes dynamically
- 📊 Monitoring: nodetool, metrics, health checks
- 💡 Best practices: Production-ready operations
💡 Key Takeaways:
- ✅ Minimum 3 nodes for production
- ✅ RF=3 + QUORUM = strong consistency
- ✅ Always decommission before removing nodes
- ✅ Run regular repairs (weekly)
- ✅ Monitor everything!
- ✅ Plan for growth proactively
🚀 Next Steps:
- ⚙️ Advanced Cluster Operations
- 📖 Data Modeling for Scale
- 🏭 Production Operations
- 🔒 Security Best Practices
🌐 You're now ready to run production Cassandra clusters!
Responsive Ad