Scaling Citus without limits: announcing Query Routers in StackGres

Citus already scaled the workers horizontally. Now the coordinators too

I’ve heard more than once that Citus is not a good sharding technology/architecture because it’s limited to one coordinator. That’s not completely untrue:

  • A single coordinator can scale clusters to very large workloads with multiple workers. But the coordinator is the sole entry point and, at some point, it can become the bottleneck. It’s a ceiling.

  • Citus supports querying workers directly since version 11, where workers behave as coordinators from a query routing/execution perspective. However, using workers as coordinators makes query routing and data processing contend for the same resources, which is far from ideal.

This post introduces the concept of Query Routers for Citus, an innovation that shipped in StackGres 1.19. StackGres already had a deep integration with Citus, making it easier than ever to create sharded clusters with Citus. Now it also adds the capability to scale coordinators horizontally; we call this Query Routers. And they remove the scaling limit that single-coordinator Citus has. With query routers there isn’t, in principle, any limit to scaling Citus.

The single coordinator limit

So where’s the limit, after all? Let’s first analyze the architecture. For standard Citus clusters, there was a single place where clients connect: the coordinator. For each query, the coordinator accepts the connection, parses the statement, plans it, works out which shard(s) hold the data to answer the query, opens or reuses connection(s) to the worker(s) that store the involved shards, forwards the query, and aggregates/relays/post-processes the result. All of that happens on one node.

Adding workers spreads the data and the storage work over more machines; the entry path still runs through one Postgres instance on one set of vCPUs, and the workers behind it can only be as busy as that one node allows.

This is the architecture of a typical (single-coordinator) Citus cluster:

Citus architecture with a single coordinator

We measured the coordinator’s limit. You can see the benchmark results at the end of the blog post.

Introducing Query Routers for Citus

Query Routers are our term for additional entry points that hold the Citus metadata, plan and route queries, store no sharded data, and are created declaratively and kept in the Citus topology by StackGres. They are optional, and you can add as many as you need to scale your Citus cluster’s entry point. The Query Routers are fronted by a Kubernetes Service, so that by pointing your database clients to the Service, you spread the load across all Query Routers.

To be clear, this is not a patch or fork of Citus. It’s more of an architectural pattern that our team designed based on our experience operating and scaling Postgres. In a nutshell, a Query Router is a single-instance worker (that can route queries as any worker since Citus 11), configured to store no data (no shards), backed by the same primitives StackGres uses to create Postgres clusters, and fronted by a Service that balances the load across all query routers (which therefore don’t need to have replicas for HA).

This pattern can actually be generalized. The coordinator still exists, and must exist, since it’s the only one that can perform DDL operations. But by using Query Routers (even if it’s one, which you can scale horizontally later if needed) its role can be limited just to DDL operations, which allows a small instance, while query routing is left to the Query Routers.

Citus architecture with query routers

Citus itself has had the building block for a long time. Any node that holds the Citus metadata (the catalog tables pg_dist_node, pg_dist_shard and pg_dist_placement) can plan a query and route it to the right worker. Since Citus 11 every worker holds that metadata, and an application can run distributed queries from any of them. A node can also be registered with its shouldhaveshards property set to false: it holds the metadata, plans and routes, and the rebalancer never places a shard on it.

See the Query Routers documentation.

In practice, a Query Router (QR) is four things at once:

  1. An SGCluster. Each QR is a single-instance SGCluster created and managed by StackGres, alongside the coordinator and the workers. Routers inherit their spec from the coordinator (except instances, which is fixed to 1), so without further configuration they run the same pod settings and the same configuration objects as the coordinator. HA in the QR layer is achieved by adding more than one QR.

  2. A Citus node in its own group range. Query routers register in pg_dist_node with a group identifier starting at 1024. Worker groups start at 1, so adding and removing routers never collides with the worker numbering (unless you plan to use more than 1023 workers; and if that’s the case, you can still raise the offset with spec.coordinator.queryRouterIndexOffset).

  3. A metadata node without shards. With Citus, data is partitioned across shards. And shards are distributed across the workers. For QRs, the coordinator registers each router’s group as an inactive placeholder with shouldhaveshards = false before the router’s Patroni is allowed to start, so no shard of a distributed table can ever be placed on it. Once the router is reachable, the coordinator activates it with citus_activate_node(). From then on the node reads isactive = true, hasmetadata = true and shouldhaveshards = false, and it holds a synchronized copy of the distributed catalog.

  4. A holder of the reference tables. Reference tables are the small tables every shard joins against, replicated to every node. A router that holds them answers queries on them locally and plans joins with them without an extra hop. A router that exists when create_reference_table runs receives the table then. For a router added later, StackGres' spec.configurations.citus.autoReplicateReferenceTables: true property lets the coordinator copy the existing reference tables to it with replicate_reference_tables('block_writes'); writes to the reference tables are blocked for the duration of that copy.

The coordinator keeps its role. DDL and metadata changes run there: create_distributed_table(), create_reference_table(), the rebalancer, adding and removing nodes, and the registration and activation of every router. Once routers exist, no application needs to connect to it for queries (unless you want to).

More about replicate_reference_tables: when used, the reference tables are copied ASAP to all the nodes (workers and QRs) and, when used with block_writes (as StackGres does), the incoming write traffic that uses such reference tables may be penalized (tip: don’t use large reference tables). For this reason autoReplicateReferenceTables is false by default.

How to use Query Routers in StackGres

Adding query routers to a Citus SGShardedCluster can be as simple as adding the queryRouterClusters property:

one coordinator
apiVersion: stackgres.io/v1beta1
kind: SGShardedCluster
metadata:
name: shop
spec:
type: citus
database: shop
postgres:
version: "18"
coordinator:
instances: 1
pods:
persistentVolume:
size: 50Gi
workers:
clusters: 3
instancesPerCluster: 2
pods:
persistentVolume:
size: 3Ti
with two query routers
apiVersion: stackgres.io/v1beta1
kind: SGShardedCluster
metadata:
name: shop
spec:
type: citus
database: shop
postgres:
version: "18"
coordinator:
instances: 1
queryRouterClusters: 2
pods:
persistentVolume:
size: 50Gi
workers:
clusters: 3
instancesPerCluster: 2
pods:
persistentVolume:
size: 3Ti

By default, Query Routers will inherit the configuration of the coordinator. But you can fully customize their configuration, instance profile and even init scripts via overrides.

Do you need help scaling your Postgres sharded cluster? Contact us.

It’s benchmark time!

To measure and profile the single coordinator’s performance limit, and to show how query routers overcome it and show near-linear scaling, we ran a benchmark based on pgbench, using the latest version of StackGres as of today, 1.19.3 (with PostgreSQL 18.4 and Citus 14.1.0).

Workers were r8gd.16xlarge (64 vCPU, 512 GiB, local NVMe), each holding about 1 TiB of pgbench_accounts in 3 shards of about 350 GiB. The coordinator was a c8g.16xlarge, 64 vCPU, as large as a worker. Query routers were the same shape. Three routers carried the same 192 vCPU as the three workers.

Queries used prepared statements, every transaction a single-row primary-key SELECT: 80% to a hot 5% of the rows, 10% uniform over all rows, and 10% to two reference tables, pgbench_tellers and pgbench_branches, 5% each. The data is twice the worker’s RAM and eight times its shared_buffers, so the uniform 10% is served from NVMe. One transaction is one query, so tps and queries per second are the same number.

Below are the graphs showing the transactions per second for the single coordinator, two and three query router cases; and their p50 and p99 latencies:

Benchmark: coordinator vs 3 query routers. TPS

Benchmark: coordinator vs 3 query routers. Latency

And the CPU usage of the coordinator, workers and query routers for the three scenarios:

Benchmark: coordinator vs 3 query routers. CPU usage

As can be seen, a single coordinator’s performance plateaus at around 550K tps with 400 connections. With all connections doing very active work (very fast queries), the coordinator maxes out CPU once the number of connections reaches the usual multiple of the core count (here: 6.25x). Adding more connections here only degrades the latency. Note that connection pooling (on by default in the coordinator on a StackGres sharded cluster) was disabled for this benchmark; it would not have helped, since every connection is active all the time. Worker nodes are not saturated and can deliver more work (if requested).

However, once we add query routers, we see close to linear scaling in performance, while keeping latency very low. In particular, with three query routers performance climbed to 1.43M tps with a p50 of 1 ms and p99 of 4 ms. And since workers were still not saturated, adding yet another query router and keeping the number of connections per query router a bit lower would have yielded even higher tps while keeping latency lower.

All in all, Query Routers, an innovation introduced in StackGres 1.19, allow Citus clusters to scale without limits. StackGres is Open Source, including the Citus integration and Query Routers.

Get started with StackGres today.