Backend Development › NoSQL & Other Data Stores
NewSQL / Distributed SQL
SQL databases that scale horizontally, like CockroachDB or Spanner.
Also known as: distributed sql, newsql, distributed sql database
Distributed SQL (sometimes “NewSQL”) is a class of databases that keep the relational model and SQL — transactions, joins, ACID — while spreading data across many machines. They aim to provide the horizontal scale of NoSQL with the consistency and query flexibility of a classic relational database: examples include CockroachDB, TiDB, YugabyteDB and Google Spanner.
The trick is that they shard data automatically and use a consensus protocol (like Raft) to replicate each shard consistently across nodes, so the cluster as a whole offers ACID transactions across the whole dataset — not just per node.
classic relational: one big machine (scale up)
NoSQL: many machines, weaker consistency (scale out)
distributed SQL: many machines, still ACID (scale out + consistency)
The classic mistakes:
- Assuming it’s free of trade-offs. Consistency and scale across a network cost latency: a write usually needs agreement among replicas, so cross-region writes are slower than a single-node commit. Distributed SQL buys scale and correctness with latency.
- Ignoring the latency of cross-shard transactions. Transactions touching data in different shards/regions need coordination; that’s where the cost concentrates. Design to keep related data co-located.
- Treating it as a drop-in for any relational database. Query planners, index choices and operational model differ; performance characteristics aren’t identical to a monolithic database. Benchmark your workload.
- Over-splitting data. A very chatty access pattern across shards defeats the design. Partitioning choices matter.
- Believing “scale out” means you can ignore schema. A good schema and index design still matter — arguably more, given distributed query plans.
- Forgetting the CAP context. Distributed SQL typically prioritises consistency (CP); during a network partition it may sacrifice availability rather than return stale data. Know the guarantee you’re getting (see CAP theorem).
When to use it: when you’ve outgrown a single-node relational database but still need SQL, transactions and strong consistency — a growing SaaS with global users, a system that can’t adopt eventual consistency. When a single node (or read replicas) still suffice, a classic relational database is simpler and lower-latency; the distributed option earns its complexity when scale or geography genuinely demands it. See sharding and partitioning.