What is NewSQL? It’s the name for databases that scale across many servers and keep SQL and real transactions. The label is old, but the idea grew up. Today the idea mostly goes by the name distributed SQL.

The Problem It Was Built For
A classic relational server gives you SQL and transactions. A transaction is a group of changes that succeed or fail together. The catch is that every write goes to one machine. When that machine is full, you buy a bigger one.
Many early NoSQL systems took the other road. They spread data across lots of cheap servers, and many gave up joins or multi-row transactions to do it. An industry analyst coined the term NewSQL in 2011. It named databases that wanted both: scale across servers and the safety of ACID.
ACID stands for atomicity, consistency, isolation and durability. In plain words, a change is all or nothing and follows the rules. It doesn’t trip over other users, and it survives a crash.
Three Ideas Behind Many Distributed SQL Databases
Under the hood, what is NewSQL? In one common design, it’s three ideas working together. The first is sharding. The database cuts a table into ranges of the key. Each range is a piece that can sit on any server. The second is replication. Each range is usually stored as three copies on three different servers, so one lost server loses nothing.
The third idea keeps the copies in agreement. One copy is the leader, and the other two follow. A write counts as committed when a majority of the copies have stored it. With three copies, a majority is two. This voting is done by a consensus protocol, such as Raft or Paxos.
One Small Example: An Orders Table
The walk-through below follows one common design: range sharding with consensus replication, as in Spanner and CockroachDB. Other systems shard by hash, and their replication and transaction details differ. Picture a bookstore with 3,000 orders. The table is keyed by OrderID. The key space splits into three ranges: orders 1 to 1000, 1001 to 2000, and 2001 to 3000. Three servers, called node 1, node 2 and node 3, hold the data.
Each range keeps one copy on every node. The leaders are spread out. Node 1 leads range 1, node 2 leads range 2, and node 3 leads range 3. That way no single node does all the write work. The picture shows the layout and one write.
How One Write Commits
The app sends UPDATE orders SET Qty = 2 WHERE OrderID = 1500; to the database. A SQL layer sits on top of every node, so the app talks to it like any SQL database. The layer looks up the key and finds that 1500 belongs to range 2, whose leader is on node 2.
The leader logs the change and sends it to the copies on node 1 and node 3. As soon as one of them confirms, two of three copies hold the change. That’s a majority, so the write commits and the app gets its answer. The slower copy catches up a moment later.
Now suppose node 2 fails. The two remaining copies of range 2 vote for a new leader, and writes continue. Nothing that was committed is lost, because every committed write already sits on at least two copies.
In CAP terms, these databases choose consistency. A range that loses its majority stops taking writes until the majority returns. I explain that choice in What Is NoSQL? Data Models, the CAP Theorem and Trade-Offs.
When One Transaction Touches Two Ranges
Real orders also change stock. Say changing order 1500 also lowers a book’s stock count. The stock row lives in another table, in a range whose leader is node 3. The transaction now spans two ranges on two leaders. Both changes must commit, or neither.
The database runs a commit protocol across the ranges, a form of two-phase commit. Each range first promises it can commit, and then all of them commit together. This works, but it costs extra network trips. Good designs keep rows that change together inside one range.
Why Not Shard by Hand?
Teams have always been able to shard on their own. The app picks a server by customer or by key, and each server runs ordinary SQL. That works, until a query needs rows from two servers. Joins across servers, transactions across servers and moving a shard when one fills up all become the app’s job.
A distributed SQL database takes that job into the database. Its SQL layer finds the right ranges and splits a big query into parts. The database also moves ranges between nodes when one node runs hot. The app still sees one database, and the code stays plain SQL.
When One Server Is Enough
The strongest objection is also the simplest: one big SQL Server is enough, and for most workloads it is. A modern server holds terabytes and handles heavy traffic. Availability groups keep copies of a database too, but all writes still go to the one primary. I covered that role in Relational Database in Big Data: Still at the Center.
Distributed SQL pays off when writes outgrow one machine. It also pays off when data must survive the loss of a data center. Google Spanner and CockroachDB are two well-known examples. The price is time. Every commit waits for a majority over the network, so a write is slower than on one local server.
What to Remember
If someone asks “What is NewSQL?”, lead with the promise: scale out and keep SQL and transactions. In the common design, three words explain how: ranges, replicas and a majority. Ranges spread the data, replicas protect it, and the majority vote keeps every copy honest. The SQL on top looks the same as always.
A design review should ask where the rows that change together live. If they sit in one range, commits stay fast. If they’re spread across ranges, every transaction pays for the trip.
NewSQL is not a product you buy, it is a promise to scale out and keep your transactions.
Published by Pinal Dave on SQLAuthority. More of my work at pinaldave.com.
Discover more from SQL Authority with Pinal Dave
Subscribe to get the latest posts sent to your email.






3 Comments. Leave new
Good One!
Nice Again Pinal
Good One