YugabyteDB is a distributed SQL database that speaks PostgreSQL. Applications connect with ordinary PostgreSQL drivers to its YSQL API and get tables, indexes, joins, foreign keys and ACID transactions, while the data underneath is split into shards, replicated with Raft across zones or regions, and rebalanced automatically when nodes join or fail. It also offers YCQL, a Cassandra-compatible API, on the same storage.

The PostgreSQL surface makes YugabyteDB easy to adopt and easy to misuse: a schema and query pattern that is fine on one Postgres server can produce hot shards, cross-region round trips or transaction retries when every row lives on a different machine. This article explains the architecture from the storage layer up, with the sharding, replication, time and transaction mechanics that decide performance, then works through a three-region example and the failure modes to plan for. Raft itself is explained in Raft log replication; here we only need what it guarantees.

Two layers and two kinds of server

Applicationsmart driver, port 5433Query layer (in every YB-TServer)YSQL (PostgreSQL) | YCQLSQLYB-Master (Raft x3)catalog, placement, balancingmetadataDocDB: each table is split into tablets; each tablet is its own Raft groupNode A (zone a)T1 leaderT2 followerT3 followerNode B (zone b)T1 followerT2 leaderT3 followerNode C (zone c)T1 followerT2 followerT3 leaderwrite to T1 leaderRaft appendA write commits when the tablet leader and one follower (majority of 3) have it in their Raft logsLeaders are spread across nodes, so every node serves writes; a lost node costs only the leaders it heldStorage per tablet: RocksDB-based regular DB + intents DB for uncommitted transactional writes
YugabyteDB layers: the query layer runs in every tablet server, YB-Master holds metadata, and DocDB stores each tablet as a Raft group whose leaders are spread across nodes.

A cluster (YugabyteDB calls it a universe) has two kinds of server. YB-Master processes, normally three, form their own Raft group and own the system catalog, tablet placement and load balancing; they are not on the data path for ordinary reads and writes. YB-TServer processes hold the data and run the query layer. YSQL reuses PostgreSQL's parser, planner and executor code, which is why compatibility is high; Yugabyte released PostgreSQL 15 compatibility in version 2.25 as a tech preview in January 2025, so check which PostgreSQL major your release series is built on before relying on a newer feature.

Below the query layer is DocDB. Every table and index is split into tablets. Each tablet is replicated, by default three times, and the replicas form an independent Raft group with one leader. Each replica stores its data in a storage engine derived from RocksDB, a log-structured merge tree, with a separate intents store for uncommitted transactional writes.

Sharding: hash, range and tablets

A row's primary key decides its tablet. YSQL supports two schemes. Hash sharding, the default for the first primary-key column, hashes the key into a 16-bit space (0 to 65535) split into ranges, one per tablet; yb_hash_code() returns that value for a key. It spreads writes evenly but makes range scans on the hashed column touch every tablet. Range sharding, declared with ASC or DESC, keeps keys in order so range scans touch few tablets, but monotonically increasing keys such as timestamps or sequences send every insert to the last tablet.

-- Hash on customer, order within a customer by recency: point lookups and
-- "latest orders for customer X" both hit one tablet.
CREATE TABLE orders (
    customer_id bigint,
    order_id    bigint,
    created_at  timestamptz NOT NULL DEFAULT now(),
    status      text,
    total       numeric(12,2),
    PRIMARY KEY ((customer_id) HASH, order_id DESC)
) SPLIT INTO 12 TABLETS;

-- Range-sharded event log: good for time-range scans, but every insert lands on the
-- newest tablet. Prefix with a small bucket to spread the hot spot.
CREATE TABLE events (
    bucket smallint,           -- e.g. event_id % 8, computed by the app
    ts     timestamptz,
    event_id uuid,
    body   jsonb,
    PRIMARY KEY ((bucket) HASH, ts ASC, event_id)
);

Tablets also split automatically as they grow, so SPLIT INTO mainly sets a sensible starting point for a table you know will be busy. For broader background on choosing shard keys see sharding strategies.

The write and read path

A write to one row goes to the leader of that row's tablet. The leader appends it to its Raft log, replicates to the followers, and acknowledges once a majority (two of three) has it durably. With replicas in three zones, losing one zone loses no committed data, and the surviving replicas elect new leaders for the tablets whose leaders were in the lost zone, typically within a few seconds.

Reads go to the tablet leader by default and are strongly consistent without a Raft round trip, because leaders hold leader leases: a newly elected leader waits until the old leader's lease has expired before serving, so two leaders never both serve reads. That gives single-row reads and writes the latency of one hop to the leader plus, for writes, one replication round trip to the nearest follower.

Hybrid time and distributed transactions

Multi-row transactions need a global order without a global clock. YugabyteDB uses hybrid time, a hybrid logical clock that combines the physical clock with a logical counter and advances whenever a node sees a later timestamp in a message; the idea is explained in hybrid logical clocks. Correctness depends on a bound on clock skew between nodes, the max_clock_skew_usec setting, which defaults to 500 ms. When a read encounters a value written within that uncertainty window after its read time, it cannot tell whether the write happened first, so the read restarts at a later time; under contention you may see this surface as a read-restart error.

A distributed transaction works like this. The transaction gets a record in a transaction status tablet. Each write is stored as a provisional record (an intent) in the target tablet's intents store, replicated through Raft like any write. Commit is a single Raft write that flips the status record to committed with a commit time; the intents are then applied to the regular store in the background. A reader that meets an intent checks the status tablet to decide whether to see it. Transactions that touch only one tablet take a faster path without the status record, which is one more reason to design keys so related rows share a tablet.

YSQL offers Serializable, Snapshot (which is what Repeatable Read maps to) and Read Committed. Read Committed historically required the TServer flag yb_enable_read_committed_isolation; without it, Read Committed silently runs as Snapshot. New universes on v2025.2 or later deployed with yugabyted, YugabyteDB Anywhere or Aeon enable it by default, but check the flag on older or upgraded clusters because the isolation level changes how conflicts surface. The semantics of snapshot isolation itself are covered in snapshot isolation.

Retries and smart drivers

Under Snapshot and Serializable, conflicting transactions fail with SQLSTATE 40001 and must be retried by the client. Every YugabyteDB application needs a retry wrapper, and the transaction body must be safe to run twice:

import random, time
import psycopg2  # pip install psycopg2-yugabytedb: this DSN needs the smart driver

DSN = ("host=yb-1.internal port=5433 dbname=shop user=app "
       "load_balance=true topology_keys=aws.us-east-1.us-east-1a")

def run_txn(conn, fn, attempts=5):
    for i in range(attempts):
        try:
            with conn:                       # commit on success, rollback on error
                with conn.cursor() as cur:
                    return fn(cur)
        except psycopg2.Error as e:
            if e.pgcode != "40001" or i == attempts - 1:
                raise
            time.sleep((2 ** i) * 0.02 * random.uniform(0.5, 1.5))

def move_stock(cur):
    cur.execute("UPDATE stock SET qty = qty - 1 WHERE sku = %s AND qty > 0", ("A-17",))
    if cur.rowcount == 0:
        raise ValueError("out of stock")
    cur.execute("INSERT INTO reservations (sku, order_id) VALUES (%s, %s)", ("A-17", 9001))

The load_balance and topology_keys parameters are understood only by YugabyteDB's smart drivers. They spread connections across tablet servers using the cluster's own list of nodes, and the topology key keeps connections in the application's zone, so you do not need a separate load balancer in front of the database.

Placing data: follower reads, tablespaces, xCluster, colocation

Three tools control where data lives and how far reads travel. Follower reads let read-only transactions read from the nearest replica at a slightly older timestamp. Enable them per session with yb_read_from_followers (default off); yb_follower_read_staleness_ms sets the staleness, defaults to 30,000 ms, and must exceed twice the maximum clock skew. They apply only to read-only transactions, and the read is stale even when the nearest replica is the leader.

SET yb_read_from_followers = true;
SET yb_follower_read_staleness_ms = 5000;
START TRANSACTION READ ONLY;
SELECT status, total FROM orders WHERE customer_id = 42 ORDER BY order_id DESC LIMIT 10;
COMMIT;

-- Pin a table's replicas to one region for data residency.
CREATE TABLESPACE eu_only WITH (replica_placement='{"num_replicas": 3, "placement_blocks": [
  {"cloud":"aws","region":"eu-west-1","zone":"eu-west-1a","min_num_replicas":1},
  {"cloud":"aws","region":"eu-west-1","zone":"eu-west-1b","min_num_replicas":1},
  {"cloud":"aws","region":"eu-west-1","zone":"eu-west-1c","min_num_replicas":1}]}');
CREATE TABLE eu_customers (id bigint PRIMARY KEY, email text) TABLESPACE eu_only;

Tablespaces with replica_placement pin tables, indexes or partitions to regions, which combined with list partitioning gives row-level geo-partitioning; the design is covered in geo-partitioning and data residency. xCluster replication asynchronously copies data between two separate universes, for disaster recovery or region-local writes, at the price of a lag window and no cross-universe transactions. Colocation (CREATE DATABASE app WITH COLOCATION = true) stores many small tables in a single tablet, avoiding thousands of tiny Raft groups for schemas with many lookup tables.

Worked example: three regions

Worked example. An order service runs in us-east-1, us-west-2 and eu-west-1 with one replica of each tablet per region. Assume illustrative round trips of 70 ms east to west and 80 ms east to Europe, and that tablet leaders are preferred in us-east-1, where most traffic originates.

A single-row insert from an east application server goes to the east leader, which needs one more replica: the nearest is the west replica at 70 ms. Commit latency is about 70 ms plus local work. A two-tablet transaction (order row plus stock row) adds the status-tablet write and intent resolution, so expect roughly two to three replication round trips, about 150 to 220 ms. A read from a Europe application server goes to the east leader: about 80 ms. With follower reads and 5 seconds of acceptable staleness, the same read is served by the local Europe replica in a few milliseconds.

The design conclusions follow from the arithmetic: keep leaders near writers, let distant readers use follower reads where staleness is acceptable, keep each transaction to as few tablets as possible, and if Europe needs fast writes for its own customers, geo-partition those rows so their leaders and replicas live in Europe.

Failure modes

Failure modes to plan for:

  • Hot tablet. A range-sharded timestamp or sequence key sends all inserts to one leader. Hash or bucket the leading key; watch per-tablet write rates.
  • Scatter-gather queries. A range filter on a hashed column scans every tablet. Check EXPLAIN (ANALYZE, DIST) output for the number of storage requests.
  • Retry storms. Hot rows plus Serializable or Snapshot produce many 40001 errors. Retry with jittered backoff and redesign hot counters.
  • Clock problems. Correctness assumes skew stays within max_clock_skew_usec. Run chrony or a cloud time service on every node and alert on drift.
  • Too many tablets. Thousands of small tables each with several tablets waste memory and Raft heartbeats. Use colocation.
  • Assuming PostgreSQL behaviour. Some extensions and features differ; sequences are cached per connection so values have gaps. Test your schema on the target release.
  • Minority partition. A region cut off from the others cannot elect leaders and stops accepting writes for its tablets; that is the price of consistency.

Trade-offs

Compared with a single PostgreSQL server, YugabyteDB trades per-query latency (every write is a network round trip to a majority) for horizontal scale and zone or region survival without manual failover. Compared with CockroachDB, it makes a similar architectural bet, a Raft group per shard and hybrid time, but reuses PostgreSQL's query layer code and adds a Cassandra-compatible API. Choose it when you need PostgreSQL semantics with more write capacity or availability than one primary gives; stay on PostgreSQL with replicas when one machine can hold the write load and minutes of failover are acceptable.

What to do next

  1. Start a local three-node cluster with yugabyted and load a copy of your schema.
  2. For each large table, decide hash or range sharding from its dominant queries, and bucket any monotonically increasing leading key.
  3. Run your top ten queries with EXPLAIN (ANALYZE, DIST) and fix any that fan out to every tablet.
  4. Add a 40001 retry wrapper around every transaction and make transaction bodies idempotent.
  5. Check yb_enable_read_committed_isolation and decide which isolation level the application really expects.
  6. Use the smart driver with load_balance and topology_keys.
  7. For multi-region, set preferred leader regions, enable follower reads for stale-tolerant reads, and test a region failure before go-live.
Key takeaway: YugabyteDB runs a PostgreSQL-compatible query layer over DocDB, which splits each table into tablets and replicates every tablet as its own Raft group, so writes cost a majority round trip and survive a zone or region loss. Hybrid time and provisional records give distributed ACID transactions that clients must retry on 40001. Pick hash or range sharding from your queries, keep transactions on few tablets, use follower reads and tablespaces to control distance, and check the isolation level your cluster really runs.