Andrew Mercer
on this page

PostgreSQL HA with Citus

Citus turns PostgreSQL into a distributed database by sharding tables across worker nodes, coordinated by one coordinator node clients actually connect to. This is a different scaling model from Replication — replication copies the same data to multiple read-only standbys, whereas Citus splits data across nodes so no single node holds the whole table. See Overview for base PostgreSQL concepts referenced below.

Architecture

  • Coordinator — the node clients connect to; holds metadata about how data is distributed and routes/plans queries across workers
  • Workers — hold the actual shards (physical PostgreSQL tables under the hood) and execute the query fragments the coordinator routes to them

The coordinator is a single point of query entry, not a single point of data storage — this matters when planning HA for the coordinator itself (see Notes and Gaps below).

Create the Infrastructure

docker-compose.yaml

services:
  coordinator:
    image: citusdata/citus
    container_name: citus_coordinator
    environment:
      POSTGRES_PASSWORD: mypassword
    ports:
      - "5432:5432"
    networks:
      - citus
    volumes:
      - coordinator_data:/var/lib/postgresql/data

  worker1:
    image: citusdata/citus
    container_name: citus_worker1
    environment:
      POSTGRES_PASSWORD: mypassword
    networks:
      - citus
    volumes:
      - worker1_data:/var/lib/postgresql/data

  worker2:
    image: citusdata/citus
    container_name: citus_worker2
    environment:
      POSTGRES_PASSWORD: mypassword
    networks:
      - citus
    volumes:
      - worker2_data:/var/lib/postgresql/data

networks:
  citus:

volumes:
  coordinator_data:
  worker1_data:
  worker2_data:
docker-compose up -d

Only the coordinator publishes port 5432 to the host — workers are reached by the coordinator over the internal citus network, not by clients directly.

Initialize the Cluster

Connect to the coordinator and register each worker:

docker exec -it citus_coordinator psql -U postgres
SELECT master_add_node('worker1', 5432);
SELECT master_add_node('worker2', 5432);
 master_add_node
-----------------
               1
(1 row)

 master_add_node
-----------------
               2
(1 row)

Check the Status of the Cluster

SELECT * FROM master_get_active_worker_nodes();
 node_name | node_port
-----------+-----------
 worker2   |      5432
 worker1   |      5432
(2 rows)

Create a Distributed Table and Load Test Data

CREATE DATABASE citus_test;

Connect to citus_test, then:

CREATE EXTENSION IF NOT EXISTS citus;

CREATE TABLE users (
    id bigserial,
    name text
);

Distributing the table shards it across the registered workers, keyed on the chosen distribution column:

SELECT create_distributed_table('users', 'id');

Choosing the right distribution column matters more than anything else in a Citus schema design — queries that filter or join on the distribution column can be routed to a single shard efficiently, while queries that don't need to fan out across every worker and recombine results on the coordinator. id works fine for this smoke test; for a real schema, pick the column most of your queries actually filter or join on (e.g. tenant_id in a multi-tenant app).

INSERT INTO users (name) SELECT 'user_' || g FROM generate_series(1, 1000) g;
INSERT 0 1000
SELECT count(*) FROM users;
 count
-------
  1000
(1 row)

Inspect shard placement directly:

SELECT shardid, nodename, nodeport
FROM pg_dist_shard_placement
ORDER BY shardid;
 shardid | nodename | nodeport
---------+----------+----------
  102008 | worker1  |     5432
  102009 | worker2  |     5432
  102010 | worker1  |     5432
  102011 | worker2  |     5432
  ...
  102039 | worker2  |     5432
(32 rows)

By default Citus created 32 shards for this table, alternating placement across the two workers — shard count is configurable via citus.shard_count before distributing the table, if the default doesn't suit your data volume.

Check That Shards Are Properly Balanced

Run a query on every worker and aggregate the per-node row counts back on the coordinator:

SELECT nodename, count(*) AS rows_on_worker
FROM run_command_on_workers(
  $cmd$ SELECT current_setting('citus.node_name') AS nodename, count(*) FROM users $cmd$
) AS t(nodename text, count bigint)
GROUP BY nodename
ORDER BY nodename;

An even split confirms the distribution column is spreading rows well. A lopsided result — especially as data grows — usually means the distribution column has low cardinality or a skewed value distribution (a small number of very common values), which is worth catching early since re-sharding an existing large distributed table is far more disruptive than choosing well up front.

Notes and Gaps

These aren't covered by the source notes this page is built from, but are worth flagging before running Citus anywhere beyond a test:

  • Coordinator HA isn't handled by anything above — the compose file here runs a single coordinator container with no standby. In production, the coordinator itself typically needs its own physical replication standby, since losing it stops all client access even though the worker data is intact.
  • Worker HA is separate from coordinator HA — each worker is a regular PostgreSQL instance under the hood and can (should) have its own standby via the same Replication mechanisms, independent of Citus.
  • Rebalancing — as workers are added or data grows unevenly, citus_rebalance_start() (or the older rebalance_table_shards()) redistributes shards across the current worker set. Not exercised in the notes above, but this is the tool to reach for if the balance check turns up lopsided.
  • Overview — base PostgreSQL concepts this builds on
  • Replication — what to layer on top for coordinator/worker HA, since Citus itself doesn't provide it