Production-Ready Clusters

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

Node 1 Tokens: 0 - 3074 Node 2 Tokens: 3075 - 6147 Node 3 Tokens: 6148 - 9223 Replication 3-Node Cluster Data replicated across all nodes
🎯

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!

1

Create docker-compose.yml

version: '3.8' services: cassandra-node1: image: cassandra:4.1 container_name: cassandra-node1 ports: - "9042:9042" # CQL port - "7199:7199" # JMX port environment: - CASSANDRA_CLUSTER_NAME=ProdCluster - CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch - CASSANDRA_DC=datacenter1 - CASSANDRA_RACK=rack1 volumes: - cassandra-node1-data:/var/lib/cassandra networks: - cassandra-net cassandra-node2: image: cassandra:4.1 container_name: cassandra-node2 environment: - CASSANDRA_CLUSTER_NAME=ProdCluster - CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch - CASSANDRA_DC=datacenter1 - CASSANDRA_RACK=rack1 - CASSANDRA_SEEDS=cassandra-node1 volumes: - cassandra-node2-data:/var/lib/cassandra networks: - cassandra-net depends_on: - cassandra-node1 cassandra-node3: image: cassandra:4.1 container_name: cassandra-node3 environment: - CASSANDRA_CLUSTER_NAME=ProdCluster - CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch - CASSANDRA_DC=datacenter1 - CASSANDRA_RACK=rack1 - CASSANDRA_SEEDS=cassandra-node1 volumes: - cassandra-node3-data:/var/lib/cassandra networks: - cassandra-net depends_on: - cassandra-node1 networks: cassandra-net: driver: bridge volumes: cassandra-node1-data: cassandra-node2-data: cassandra-node3-data:
2

Start the Cluster

# Start all nodes docker-compose up -d # Wait 2-3 minutes for cluster to form # Nodes need time to discover each other via gossip # Watch logs docker-compose logs -f
3

Verify Cluster Status

# Check cluster status docker exec cassandra-node1 nodetool status # Expected output: 3 nodes UN (Up and Normal) Datacenter: datacenter1 ======================= Status=Up/Down |/ State=Normal/Leaving/Joining/Moving -- Address Load Tokens Owns Host ID UN 172.18.0.2 69 KiB 16 33.3% abc123... UN 172.18.0.3 65 KiB 16 33.3% def456... UN 172.18.0.4 73 KiB 16 33.4% ghi789... # All nodes show UN = SUCCESS! ✅

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!

1

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)
2

Create Keyspace with RF=3

-- Connect to any node docker exec -it cassandra-node1 cqlsh -- Create keyspace with replication CREATE KEYSPACE prod_app WITH REPLICATION = { 'class': 'SimpleStrategy', 'replication_factor': 3 }; -- Use the keyspace USE prod_app; -- Create table CREATE TABLE users ( id UUID PRIMARY KEY, name TEXT, email TEXT, created_at TIMESTAMP );
3

Test Replication

-- Insert data on node 1 INSERT INTO users (id, name, email, created_at) VALUES (uuid(), 'Alice', 'alice@example.com', toTimestamp(now())); INSERT INTO users (id, name, email, created_at) VALUES (uuid(), 'Bob', 'bob@example.com', toTimestamp(now())); exit # Connect to node 2 - data is there! docker exec -it cassandra-node2 cqlsh USE prod_app; SELECT * FROM users; -- You'll see the data! ✅ -- This proves replication is working!
✅

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

-- Set consistency level in cqlsh CONSISTENCY QUORUM; -- Now all queries use QUORUM SELECT * FROM users; INSERT INTO users (...) VALUES (...); -- Change to ONE for this session CONSISTENCY ONE; -- Check current level CONSISTENCY;

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

# Add to docker-compose.yml cassandra-node4: image: cassandra:4.1 container_name: cassandra-node4 environment: - CASSANDRA_CLUSTER_NAME=ProdCluster - CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch - CASSANDRA_DC=datacenter1 - CASSANDRA_SEEDS=cassandra-node1 volumes: - cassandra-node4-data:/var/lib/cassandra networks: - cassandra-net # Start the new node docker-compose up -d cassandra-node4 # Wait 5-10 minutes for bootstrapping # New node will stream data from existing nodes # Check status docker exec cassandra-node1 nodetool status # Should show 4 nodes UN
🔄

Run Cleanup After Adding Nodes

# After new node joins, cleanup old nodes # This removes data they no longer own docker exec cassandra-node1 nodetool cleanup docker exec cassandra-node2 nodetool cleanup docker exec cassandra-node3 nodetool cleanup # Frees up disk space!
➖

Removing a Node (Decommission)

# Gracefully remove node (streams data to others) docker exec cassandra-node4 nodetool decommission # Wait for decommission to complete (5-30 min) # Node will stream its data to other nodes # Check status from another node docker exec cassandra-node1 nodetool status # Node 4 should be gone # Now stop the container docker stop cassandra-node4 docker rm cassandra-node4

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

# View all nodes nodetool status # Expected: All nodes UN # View cluster info nodetool describecluster # View ring token distribution nodetool ring

Monitor Data Distribution

# Check disk usage per node nodetool status # Look at "Load" column # Should be roughly equal across nodes # Example: 100GB, 98GB, 102GB ✅ # Bad: 200GB, 50GB, 50GB ❌

Monitor Performance

# Table statistics nodetool tablestats keyspace_name.table_name # Shows: # - Read/write count # - Read/write latency # - Disk space used # Thread pool stats nodetool tpstats # Shows pending/active tasks

Check Replication

# Verify data is replicated # Insert on node 1, read from node 2 # If you can read from any node: ✅ # If data missing on some nodes: ❌ run repair # Repair to sync replicas nodetool repair

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

-- Even for single datacenter! CREATE KEYSPACE prod_ks WITH REPLICATION = { 'class': 'NetworkTopologyStrategy', 'datacenter1': 3 }; -- Easier to add datacenters later

✅ Run Regular Repairs

# Schedule weekly repairs # Ensures replicas stay in sync # Full repair (once per week) nodetool repair -full # Or use Reaper tool for automated repairs

✅ Take Regular Snapshots

# Daily snapshots on ALL nodes nodetool snapshot -t backup-$(date +%Y%m%d) # Store snapshots offsite (S3, etc) # Test restores regularly!

✅ 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:

🌐 You're now ready to run production Cassandra clusters!

Advertisement

Responsive Ad