Horizontal Database Sharding at Scale: Citus Distributed Tables, Distributed Transactions, and Partition-Wise Joins

When a single PostgreSQL primary reaches write saturation and storage limits, Citus transforms PostgreSQL into a distributed cluster. Master shard keys, 2PC distributed transactions, and co-located joins.

The Limits of Vertical Scaling: When Read Replicas and Hardware Upgrades Fail

In high-growth SaaS applications and transactional platforms, PostgreSQL begins life as a single monolithic instance. As read traffic climbs, teams scale out by provisioning read replicas behind PgBouncer or connection poolers. But as data volumes scale into multiple terabytes and write throughput reaches tens of thousands of IOPS, the monolithic architecture hits an unyielding wall:

  • Write Bottlenecks: Read replicas cannot absorb write transactions; all INSERT, UPDATE, and DELETE operations must hit the single primary node.
  • Disk I/O and Buffer Pool Saturation: When the active working set (tables and indexes) surpasses physical RAM, cache hit ratios plummet, forcing PostgreSQL to read pages continuously from NVMe storage.
  • Autovacuum Freezing and Table Bloat: Vacuuming multi-terabyte tables requires immense I/O bandwidth. Checkpointing and transaction ID wraparound vacuums begin impacting client latency.

When the largest available cloud instance (e.g., 128 vCPUs and 1TB of RAM) is insufficient, the system must scale horizontally. Rather than abandoning PostgreSQL for a NoSQL store and forfeiting ACID transactions, relational joins, and rich indexing, teams can leverage Citus to transform PostgreSQL into a horizontally distributed database cluster.

The Citus Distributed Architecture: Coordinator and Worker Nodes

Citus extends PostgreSQL from within using open-source extension hooks without forking the codebase. A Citus cluster consists of two node types:

  1. Coordinator Node: Acts as the single entry point for client applications. The coordinator stores cluster metadata (which shards reside on which worker nodes), parses incoming SQL queries, generates distributed execution plans, pushes optimized query fragments out to worker nodes, and aggregates final result sets.
  2. Worker Nodes: Standalone PostgreSQL instances that hold physical shards (standard PostgreSQL tables containing slices of distributed data). Workers execute queries in parallel against their local data, utilizing their own local CPU cores, RAM, and NVMe drives.

Clients connect to the Citus coordinator using standard PostgreSQL drivers (such as psycopg3, asyncpg, or JDBC), completely unaware that data is distributed across dozens of physical machines.

The Art of the Distribution Column and Shard Co-Location

The single most critical design decision in a distributed database is selecting the Distribution Column (Shard Key). In multi-tenant B2B platforms, the optimal distribution column is almost universally the tenant identifier (e.g., company_id or tenant_id). In event or analytics platforms, it might be user_id or device_id.

The Principle of Co-Location

When multiple related tables are distributed on the same distribution column with matching shard counts, Citus automatically co-locates their shards:

-- Enable Citus extension
CREATE EXTENSION IF NOT EXISTS citus;

-- 1. Create multi-tenant schema with company_id composite primary keys
CREATE TABLE companies (
    id BIGSERIAL PRIMARY KEY,
    name TEXT NOT NULL,
    created_at TIMESTAMPTZ DEFAULT NOW()
);

CREATE TABLE users (
    id BIGSERIAL,
    company_id BIGINT REFERENCES companies(id),
    email TEXT NOT NULL,
    PRIMARY KEY (company_id, id)
);

CREATE TABLE orders (
    id BIGSERIAL,
    company_id BIGINT REFERENCES companies(id),
    user_id BIGINT NOT NULL,
    total_amount NUMERIC(12, 2) NOT NULL,
    PRIMARY KEY (company_id, id)
);

-- 2. Distribute tables using company_id as the shard key
SELECT create_distributed_table('companies', 'id');
SELECT create_distributed_table('users', 'company_id');
SELECT create_distributed_table('orders', 'company_id');

Because all records belonging to company_id = 42 across companies, users, and orders reside on the exact same physical worker node, complex multi-table joins execute entirely locally on that worker. No network shuffles, no cross-node data transfers, and no distributed locking overhead.

Distributed Query Execution: Partition-Wise Joins & 2-Phase Commits

When a client executes a query filtering by the shard key:

SELECT u.email, SUM(o.total_amount) 
FROM users u 
JOIN orders o ON u.company_id = o.company_id AND u.id = o.user_id 
WHERE u.company_id = 42 
GROUP BY u.email;

The coordinator performs Router Query Planning: it inspects the metadata cache, identifies the single worker node hosting the shards for company_id = 42, and routes the SQL query directly to that worker. The query runs at native, bare-metal PostgreSQL speed.

Cross-Shard Analytics via Partition-Wise Joins

If an analytical query spans across all tenants without filtering on company_id, Citus executes a Distributed Parallel Query Plan:

SELECT company_id, COUNT(*) FROM orders GROUP BY company_id;

The coordinator divides the query into sub-queries corresponding to each shard, sends them asynchronously across all worker nodes in parallel, and merges the partial aggregates. A cluster with 16 worker nodes each running 16 cores harnesses 256 CPU cores and 16 NVMe disk buses simultaneously, achieving orders of magnitude faster execution than a single monolithic instance.

Distributed Transactions & ACID Durability

For write transactions mutating shards across multiple worker nodes, Citus coordinates a Two-Phase Commit (2PC) protocol (PREPARE TRANSACTION and COMMIT PREPARED) managed by an internal distributed transaction manager, ensuring full ACID guarantees and distributed deadlock detection.

Partitioning Prerequisites: Before migrating to distributed Citus tables, understand single-instance partitioning fundamentals in PostgreSQL Declarative Partitioning at 50M Rows/Day with pg_partman.

Step-by-Step Migration from Vanilla PostgreSQL to Citus

Migrating a monolithic database to a distributed Citus cluster follows a disciplined architectural roadmap:

  1. Enforce Composite Primary & Foreign Keys: PostgreSQL unique constraints and primary keys must include the distribution column. Alter your relational schemas to include tenant_id in all primary and foreign key constraints.
  2. Classify Tables into Distributed, Reference, and Local:
    • Distributed Tables: Large tenant-partitioned data (orders, events, invoices).
    • Reference Tables: Small, lookup tables (currencies, countries, role permissions). Replicated to all worker nodes via create_reference_table('countries') so they can be joined locally against any shard.
    • Local Tables: Administrative tables that remain solely on the coordinator node.
  3. Distribute Tables Concurrently: Citus supports online, non-blocking shard redistribution, enabling zero-downtime scaling as cluster capacity grows.

Horizontal database sharding via Citus bridges the gap between the familiar developer ergonomics of PostgreSQL and the boundless horizontal scalability demanded by planetary-scale applications.

All Insights
Chat on WhatsApp