Citus sharding technology

Citus Use Cases

Multi-Tenant

The multi-tenant architecture uses hierarchical database modeling to distribute queries across nodes. The tenant ID is stored in a column on each table, and Citus routes queries to the appropriate worker node.

Best practices:

  • Partition distributed tables by a common tenant_id column
  • Convert small cross-tenant tables to reference tables
  • Ensure all queries filter by tenant_id

Real-Time Analytics

Real-time architectures depend on specific distribution properties to achieve highly parallel processing.

Best practices:

  • Choose a column with high cardinality as the distribution column
  • Choose a column with even distribution to avoid skewed data
  • Distribute fact and dimension tables on their common columns

Time-Series

Important: Do NOT use the timestamp as the distribution column for time-series data. A hash distribution based on time distributes times seemingly at random, leading to network overhead for range queries.

Best practices:

  • Use a different distribution column (tenant_id or entity_id)
  • Use PostgreSQL table partitioning for time ranges

Co-located Tables

Co-located tables are distributed tables that share common columns in the distribution key. This improves performance since distributed queries avoid querying more than one Postgres instance for correlated columns.

Benefits of co-location:

  • Full SQL support for queries on a single set of co-located distributed partitions
  • Multi-statement transaction support for modifications
  • Aggregation through INSERT..SELECT
  • Foreign keys between co-located tables
  • Distributed outer joins
  • Pushdown CTEs (PostgreSQL >= 12)

Example:

SELECT create_distributed_table('event', 'tenant_id');
SELECT create_distributed_table('page', 'tenant_id', colocate_with => 'event');

Reference Tables

Reference tables are replicated across all worker nodes and automatically kept in sync during modifications. Use them for small tables that need to be joined with distributed tables.

SELECT create_reference_table('geo_ips');

Scaling Workers

Adding a new worker is simple - increase the clusters field value in the workers section:

apiVersion: stackgres.io/v1beta1
kind: SGShardedCluster
metadata:
  name: my-sharded-cluster
spec:
  workers:
    clusters: 3  # Increased from 2

After provisioning, rebalance data using the resharding operation:

apiVersion: stackgres.io/v1
kind: SGShardedDbOps
metadata:
  name: reshard
spec:
  sgShardedCluster: my-sharded-cluster
  op: resharding
  resharding:
    citus: {}

Query Routers

Query routers are special workers that route read/write queries to the workers that store the sharded data, but do not store sharded table data themselves. Like the coordinator, a query router is an SGCluster that participates in the Citus topology, but its shouldhaveshards flag is set to false so the rebalancer never assigns shards to it.

Query routers allow you to scale the number of read/write entrypoints horizontally without overloading the workers that store the sharded table data. This is useful when the workload is bottlenecked on the coordinator’s connection handling, the parser/planner, or the cross-shard query coordination, rather than on the workers' I/O.

Query routers are only supported when spec.type is citus.

Architecture

Each query router is a single-instance SGCluster created and managed by the operator alongside the coordinator and the workers. Query routers inherit their spec from the coordinator (apart from instances, which is fixed to one) so they share the same Postgres version, configuration, pooling and pods configuration. They participate in the Citus topology as additional Citus nodes, with a dedicated groupid derived from the query router index offset (see Index Offset).

Query routers do not have replicas: each query router cluster is composed of a single Pod. To scale the number of query router entrypoints, increase coordinator.queryRouterClusters.

Enabling Query Routers

Set spec.coordinator.queryRouterClusters to the desired number of query router workers:

apiVersion: stackgres.io/v1beta1
kind: SGShardedCluster
metadata:
  name: cluster
spec:
  type: citus
  database: mydatabase
  coordinator:
    instances: 2
    queryRouterClusters: 2
    pods:
      persistentVolume:
        size: '10Gi'
  workers:
    clusters: 4
    instancesPerCluster: 2
    pods:
      persistentVolume:
        size: '10Gi'

The example above creates two query router SGClusters (one Pod each) in addition to the coordinator and the four workers.

Naming

By default the query router SGClusters are named after the SGShardedCluster name with the -router suffix and the zero-based index appended. For example, with an SGShardedCluster called cluster and queryRouterClusters: 2, the operator creates the SGClusters cluster-router0 and cluster-router1.

You can override the template via spec.coordinator.queryRouterClusterNameTemplate:

spec:
  coordinator:
    queryRouterClusters: 2
    queryRouterClusterNameTemplate: my-router

This produces my-router0 and my-router1 instead. As with the workers' clusterNameTemplate, this field can only be set on creation.

Index Offset

Query routers register as Citus nodes with a groupid starting at 1024. This offset keeps the regular workers' group identifiers (which start at 1) separated from the query routers' group identifiers, so adding or removing query routers never collides with the worker group numbering.

If you plan to have more than 1023 regular worker SGClusters, raise the offset accordingly with spec.coordinator.queryRouterIndexOffset (minimum 1024):

spec:
  coordinator:
    queryRouterIndexOffset: 4096
    queryRouterClusters: 2

Connecting Through Query Routers

Each query router SGCluster exposes a Kubernetes Service named after the SGCluster (for example cluster-router0). Applications can connect through any of these services to issue read/write queries; the query router will forward the query to the appropriate worker.

You can enable, disable or customize the type of the query router primary services centrally through spec.postgresServices.coordinator.queryRouters:

spec:
  postgresServices:
    coordinator:
      queryRouters:
        type: LoadBalancer
        enabled: true

This setting is propagated to the primary Service of every query router SGCluster. Replicas Services on query router SGClusters are always disabled because query routers are single-instance clusters.

Overriding Specific Query Routers

Like worker overrides, you can override individual query router clusters via spec.workers.overrides by setting type: QueryRouter on the override entry. The index (or indexes) refers to the zero-based query router identifier (i.e. 0 selects the first query router), not the offset Citus group identifier.

Whatever an override entry does not set is inherited from spec.coordinator, as the rest of the query router spec. A query router with no override therefore uses spec.coordinator.sgInstanceProfile and spec.coordinator.configurations (sgPostgresConfig and sgPoolingConfig); the spec.workers ones are never applied to a query router. Since the index spaces of the two types are distinct, a type: Worker entry and a type: QueryRouter entry may share the same index.

Give the query routers their own Postgres configuration by referencing it from a type: QueryRouter entry:

spec:
  coordinator:
    queryRouterClusters: 2
  workers:
    clusters: 4
    overrides:
    - indexes: ["all"]
      type: QueryRouter
      configurations:
        sgPostgresConfig: routers-postgres-config

For citus sharded clusters the operator does not use the referenced SGPostgresConfig directly: it generates a per query router SGPostgresConfig named <query router SGCluster name>-<postgres major version> out of it, adding the parameters Citus requires.

See Cluster Names and Overrides for the full reference of the index, indexes and type fields.

Scaling Query Routers

Add or remove query routers by changing spec.coordinator.queryRouterClusters:

kubectl patch sgshardedcluster my-sharded-cluster --type merge \
  -p '{"spec":{"coordinator":{"queryRouterClusters":3}}}'

Query routers do not store sharded data, so adding or removing them does not require resharding. The coordinator registers each query router in the Citus node table (pg_dist_node) with the shouldhaveshards flag set to false before Patroni of the query router is started, so that no shard of a distributed table can ever be placed on a query router. The SGCluster of a query router is generated with .spec.configurations.patroni.startGateAnnotations and the operator sets the corresponding annotation on it only once the coordinator reports the query router as registered. When .spec.replicateFrom is set the SGShardedCluster is a replica whose pg_dist_node is replicated from the source, so the SGCluster of a query router is generated without the start gate and its Patroni is started immediately.

The coordinator updates the Citus node table on a schedule that can be changed with .spec.configurations.citus.updateNodeInterval (every 10 seconds by default):

apiVersion: stackgres.io/v1beta1
kind: SGShardedCluster
metadata:
  name: my-sharded-cluster
spec:
  configurations:
    citus:
      updateNodeInterval: PT10S
      enableNodeAutoRemoval: true

The nodes of the workers or query routers removed by decreasing .spec.workers.clusters or .spec.coordinator.queryRouterClusters are not removed from the Citus node table by Patroni. Since Citus can only remove a primary node that can be reached, the SGCluster of a removed worker or query router is kept running while its group is registered in the Citus node table, and is scaled down to 0 instances only once its group has been removed from it. Set .spec.configurations.citus.enableNodeAutoRemoval to true to let the coordinator remove them (it only removes the nodes that hold no shard of a distributed table and can be reached), or remove them manually on the coordinator primary:

SELECT citus_remove_node('<nodename>', <nodeport>);

NOTE: a worker that still holds shards of a distributed table is never removed. Move its shards to the remaining workers first (for example with SELECT citus_drain_node('<nodename>', <nodeport>)).

Connection Pooling Between Nodes

By default the Citus nodes connect to each other through the connection pooler (PgBouncer) running in the Pod of each node, instead of opening their connections directly to Postgres. The coordinator keeps the Citus pg_dist_poolinfo table of the coordinator, of the workers and of the query routers updated so that Citus uses the port of PgBouncer (6432, or the Envoy entry port 7432 when Envoy is enabled) instead of the Postgres port registered in pg_dist_node. Only the port is set, so the connections follow the host that Patroni updates in pg_dist_node after a failover. Citus ignores pg_dist_poolinfo for the connections that can not go through a pooler (like the ones of the shard rebalancer).

Since the nodes connect through PgBouncer, the disableConnectionPooling fields of the coordinator, the workers, the query routers and their overrides are ignored and PgBouncer is always created. To connect directly to Postgres, set .spec.configurations.citus.connectToPooler to false:

apiVersion: stackgres.io/v1beta1
kind: SGShardedCluster
metadata:
  name: my-sharded-cluster
spec:
  configurations:
    citus:
      connectToPooler: false

When disabled, the coordinator removes the entries it created from pg_dist_poolinfo on the next update of the nodes (see .spec.configurations.citus.updateNodeInterval).

Distributed Partitioned Tables

Citus allows creating partitioned tables that are also distributed for time-series workloads. With partitioned tables, removing old historical data is fast and doesn’t generate bloat:

CREATE TABLE github_events (
  event_id bigint,
  event_type text,
  repo_id bigint,
  created_at timestamp
) PARTITION BY RANGE (created_at);

SELECT create_distributed_table('github_events', 'repo_id');

SELECT create_time_partitions(
  table_name         := 'github_events',
  partition_interval := '1 month',
  end_at             := now() + '12 months'
);

Columnar Storage

Citus supports columnar storage for distributed partitioned tables. This append-only format can greatly reduce data size and improve query performance, especially for numerical values:

CALL alter_old_partitions_set_access_method(
  'github_events',
  '2015-01-01 06:00:00' /* older_than */,
  'columnar'
);

Note: Columnar storage disallows updating and deleting rows, but you can still remove entire partitions.

Creating a basic Citus Sharded Cluster

Create the SGShardedCluster resource:

apiVersion: stackgres.io/v1beta1
kind: SGShardedCluster
metadata:
  name: cluster
spec:
  type: citus
  database: mydatabase
  postgres:
    version: '15'
  coordinator:
    instances: 2
    pods:
      persistentVolume:
        size: '10Gi'
  workers:
    clusters: 4
    instancesPerCluster: 2
    pods:
      persistentVolume:
        size: '10Gi'

This configuration will create a coordinator with 2 Pods and 4 workers with 2 Pods each.

By default the coordinator node has a synchronous replica to avoid losing any metadata that could break the sharded cluster.

The workers are where sharded data lives and have a replica in order to provide high availability to the cluster.

SG Sharded Cluster

After all the Pods are Ready you can view the topology of the newly created sharded cluster by issuing the following command:

kubectl exec -n my-cluster cluster-coord-0 -c patroni -- patronictl list
+ Citus cluster: cluster --+------------------+--------------+---------+----+-----------+
| Group | Member           | Host             | Role         | State   | TL | Lag in MB |
+-------+------------------+------------------+--------------+---------+----+-----------+
|     0 | cluster-coord-0  | 10.244.0.16:7433 | Leader       | running |  1 |           |
|     0 | cluster-coord-1  | 10.244.0.34:7433 | Sync Standby | running |  1 |         0 |
|     1 | cluster-worker0-0 | 10.244.0.19:7433 | Leader       | running |  1 |           |
|     1 | cluster-worker0-1 | 10.244.0.48:7433 | Replica      | running |  1 |         0 |
|     2 | cluster-worker1-0 | 10.244.0.20:7433 | Leader       | running |  1 |           |
|     2 | cluster-worker1-1 | 10.244.0.42:7433 | Replica      | running |  1 |         0 |
|     3 | cluster-worker2-0 | 10.244.0.22:7433 | Leader       | running |  1 |           |
|     3 | cluster-worker2-1 | 10.244.0.43:7433 | Replica      | running |  1 |         0 |
|     4 | cluster-worker3-0 | 10.244.0.27:7433 | Leader       | running |  1 |           |
|     4 | cluster-worker3-1 | 10.244.0.45:7433 | Replica      | running |  1 |         0 |
+-------+------------------+------------------+--------------+---------+----+-----------+

You may also check that they are already configured in Citus by running the following command:

$ kubectl exec -n my-cluster cluster-coord-0 -c patroni -- psql -d mydatabase -c 'SELECT * FROM pg_dist_node'
 nodeid | groupid |  nodename   | nodeport | noderack | hasmetadata | isactive | noderole | nodecluster | metadatasynced | shouldhaveshards 
--------+---------+-------------+----------+----------+-------------+----------+----------+-------------+----------------+------------------
      1 |       0 | 10.244.0.34 |     7433 | default  | t           | t        | primary  | default     | t              | f
      3 |       2 | 10.244.0.20 |     7433 | default  | t           | t        | primary  | default     | t              | t
      2 |       1 | 10.244.0.19 |     7433 | default  | t           | t        | primary  | default     | t              | t
      4 |       3 | 10.244.0.22 |     7433 | default  | t           | t        | primary  | default     | t              | t
      5 |       4 | 10.244.0.27 |     7433 | default  | t           | t        | primary  | default     | t              | t
(5 rows)

Please, take into account that the groupid column of the pg_dist_node table is the same as the Patroni Group column above. In particular, the group with identifier 0 is the coordinator group (coordinator have shouldhaveshards column set to f).

For a more complete configuration please have a look at Create Citus Sharded Cluster Section.