Wide-Column Stores
Concept
If you are building Netflix or Uber, you have millions of apps sending GPS coordinates and viewing metrics to your servers every single second.
You are dealing with Massive Write Throughput.
A standard SQL database will immediately crash under this load because every INSERT requires the database to lock the B-Tree index, traverse it, and rewrite the physical pages on the hard drive (Random I/O).
A Wide-Column Store (like Apache Cassandra or ScyllaDB) is a specialized NoSQL database designed specifically to absorb millions of writes per second without breaking a sweat, distributed across dozens of massive servers.
How It Absorbs Massive Writes (LSM Trees)
Cassandra achieves insane write speeds because it completely abandons B-Tree indexes. It uses an architecture called the LSM Tree (Log-Structured Merge-Tree).
When data arrives, Cassandra does not try to find the correct row on the hard drive to update it.
It simply appends the raw data to the very end of a sequential log file in RAM (the MemTable). Appending to the end of a log is the fastest possible operation in computer science. When the RAM gets full, it dumps the log directly to the hard drive as an immutable, frozen file (an SSTable).
Wait, if it just blindly dumps unordered logs to the hard drive, how does it read data?
Reads in Cassandra are actually slower than writes. When you request a user’s profile, Cassandra has to search through multiple frozen log files on the hard drive, merge the fragments together in RAM, figure out which data is the newest (using timestamps), and return the final row.
Cassandra trades Read speed to achieve infinite Write speed.
The “Wide-Column” Data Model
It looks like a relational table, but it isn’t.
In SQL, every row must have the exact same columns.
In Cassandra, each row is actually a discrete Hash Map. Row 1 might contain 5 columns. Row 2 might contain 50,000 columns (e.g., storing a user’s entire chronological history of events as individual columns attached to their single row ID).
High Availability (No Master Node)
Standard SQL databases use a Primary-Replica architecture. If the Primary dies, writes fail until a failover happens.
Cassandra uses a Masterless Ring Architecture.
All servers (Nodes) in the cluster are exactly equal. A Node.js backend can send an INSERT statement to any Node in the cluster. That Node will use Consistent Hashing to figure out which 3 servers actually own the data, and forward the write to them.
Because there is no Single Point of Failure, you can physically unplug 3 servers from a Cassandra cluster, and the database will continue accepting 100,000 writes per second without a single error.
Interview Questions
Q: A developer builds an IoT application using Cassandra. They want to find all temperature sensors that reported over 100 degrees today. They write SELECT * FROM sensors WHERE temp > 100. The query fails. Why?
A: Because of its distributed hash-ring architecture, Cassandra distributes data randomly across the cluster based on the Partition Key (e.g., sensor_id).
To answer a WHERE temp > 100 query, Cassandra would have to ask every single server in the 50-server cluster to do a full scan of their hard drives (Scatter-Gather), which would crash the cluster.
Cassandra physically prevents you from writing WHERE clauses on unindexed columns. In Cassandra, you must design your database schema based entirely on the queries you intend to run. You cannot do ad-hoc analytical filtering like you can in SQL.
Q: Explain the concept of “Eventual Consistency” and “Read Repair” in Cassandra.
A: If your Cassandra cluster keeps 3 copies of your data (Replication Factor = 3), and you execute a Write, you can configure it so that it only waits for 1 server to acknowledge the write before returning “Success” to the user. This makes the write blazing fast.
However, the other 2 servers are now slightly out of date. If a user instantly reads from Server 2, they get stale data (Eventual Consistency).
When Cassandra notices that Server 2 returned stale data, a background process called Read Repair instantly pushes the fresh data from Server 1 to Server 2, silently fixing the inconsistency for the next user.