PostgreSQL's scalability claims are often quoted but rarely quantified. To test the ceiling on ordinary hardware, a synthetic workload was generated and loaded into a Citus cluster running entirely on a single local desktop machine.
The result, verified by a full count:
|
1 2 3 4 5 |
test=# SELECT count(*) FROM t_data ; count --------------- 1024000000000 (1 row) |
Under Citus columnar storage, the same table reports a far smaller footprint:
|
1 2 3 4 5 6 7 8 9 10 11 12 |
test=# \d t_data Table "public.t_data" Column | Type | Collation | Nullable | Default --------+---------+-----------+----------+--------- key_id | integer | | | data | bigint | | | test=# SELECT pg_size_pretty(citus_table_size('t_data')); pg_size_pretty ---------------- 1563 GB (1 row) |
Why storage layout comes first
A real dataset of that size was not available, so the rows were generated. The layout decision, however, has to be made before any data lands on disk. Running the math on plain row storage makes the constraint obvious:
|
1 2 3 4 5 |
test=# SELECT pg_size_pretty(1024000000000::int8 * (24 + 12)); pg_size_pretty ---------------- 34 TB (1 row) |
With a tuple header of roughly 24 bytes plus 12 bytes of data per row, the table would approach 34 TB — beyond the maximum table size PostgreSQL permits with 8k blocks. That leaves three viable paths: PostgreSQL table partitioning over row storage, a single columnar table, or sharding with Citus.
For query speed, the chosen combination was sharded Citus columnar storage: one coordinator plus four nodes, all on the same desktop:
|
1 2 3 4 5 6 7 8 9 10 11 |
test=# SELECT nodeid, nodename, nodeport FROM pg_dist_node ORDER BY 1; nodeid | nodename | nodeport --------+-----------+---------- 1 | localhost | 5432 5 | localhost | 6001 6 | localhost | 6002 7 | localhost | 6003 8 | localhost | 6004 (5 rows) |
Building the distributed table
Once Citus is installed and wired up, a simple table is created and distributed on one of its columns:
|
1 2 3 4 5 6 7 8 |
test=# CREATE TABLE t_data ( key_id int, data bigint ) USING columnar; CREATE TABLE test=# SELECT create_distributed_table('t_data', 'key_id'); SELECT |
Initial data comes from generate_series. This INSERT produces 62.5 million rows in the most basic form:
|
1 2 3 4 |
test=# INSERT INTO t_data SELECT id % 100000, id FROM generate_series(1, 62500000) AS id; INSERT |
From there, the row count is doubled repeatedly until the target is reached:
|
1 2 3 |
... test=# INSERT INTO t_data SELECT * FROM t_data; ... <run more often> ... |
Querying a trillion rows
The simplest possible query is a count:
|
1 2 3 4 5 6 7 |
test=# SELECT count(*) FROM t_data; count --------------- 1024000000000 (1 row) Time: 3225188.861 ms (53:45.189) |
It completes in roughly 53 minutes. During execution, Citus dispatches work across a large number of processes spanning the shards:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 |
top - 20:16:01 up 54 min, 1 user, load average: 30.02, 14.47, 6.39 Tasks: 703 total, 34 running, 669 sleeping, 0 stopped, 0 zombie %Cpu(s): 99.8 us, 0.2 sy, 0.0 ni, 0.0 id, 0.0 wa, 0.0 hi, 0.0 si, 0.0 st MiB Mem : 128650.6 total, 50517.6 free, 20438.3 used, 72092.9 buff/cache MiB Swap: 8192.0 total, 8192.0 free, 0.0 used. 108212.4 avail Mem PID USER PR NI VIRT RES SHR S %CPU %MEM TIME+ COMMAND 17845 hs 20 0 33.6g 225728 219904 R 100.0 0.2 2:33.04 postgres 16609 hs 20 0 249712 147672 144640 R 100.0 0.1 4:05.39 postgres 16610 hs 20 0 249704 149592 145984 R 100.0 0.1 3:49.71 postgres 17833 hs 20 0 247584 127732 122112 R 100.0 0.1 2:33.29 postgres 17835 hs 20 0 247588 145580 139520 R 100.0 0.1 2:32.36 postgres 17843 hs 20 0 247640 131620 125696 R 100.0 0.1 2:32.86 postgres 17850 hs 20 0 33.6g 203300 197376 R 100.0 0.2 2:33.20 postgres 17853 hs 20 0 247636 120900 115200 R 100.0 0.1 2:32.29 postgres 17855 hs 20 0 247636 125584 119808 R 100.0 0.1 2:32.84 postgres 17857 hs 20 0 33.6g 231460 225280 R 100.0 0.2 2:33.40 postgres 16608 hs 20 0 249704 148464 145140 R 99.7 0.1 3:47.70 postgres 17840 hs 20 0 247644 122632 117248 R 99.7 0.1 2:32.76 postgres 17841 hs 20 0 247644 131788 126208 R 99.7 0.1 2:32.44 postgres 17848 hs 20 0 247636 126248 120320 R 99.7 0.1 2:32.73 postgres 17854 hs 20 0 247644 126436 121088 R 99.7 0.1 2:32.59 postgres 16611 hs 20 0 249712 147620 144072 R 99.3 0.1 3:51.26 postgres ... |
The important observation is that the workload is CPU-bound rather than I/O-bound. CPU is generally easier and cheaper to scale than I/O, so pegged CPUs here are a positive signal rather than a problem.
Grouping, and where the plan pays off
Counting instances per group is a common pattern:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 |
test=# explain SELECT key_id, count(*) FROM t_data GROUP BY 1 ORDER BY 2 DESC; QUERY PLAN ---------------------------------------------------------------------------- Sort (cost=8304.82..8554.82 rows=100000 width=12) Sort Key: remote_scan.count DESC -> Custom Scan (Citus Adaptive) (cost=0.00..0.00 rows=100000 width=12) Task Count: 32 Tasks Shown: One of 32 -> Task Node: host=localhost port=5432 dbname=test -> HashAggregate (cost=165451245.99..165451249.16 rows=317 width=12) Group Key: key_id -> Custom Scan (ColumnarScan) on t_data_106299 t_data (cost=0.00..3163471.27 rows=32457554944 width=4) Columnar Projected Columns: key_id (11 rows) Time: 287.430 ms |
Citus dispatches this query to the shards, and the plan shows each shard being scanned via ColumnarScan and aggregated locally — much like a single GROUP BY on one machine. The distribution column is key_id, so every value of key_id lives on exactly one shard. Citus only has to collect the partial results and order them as the query requests.
To see this at a lower level, connect to one of the database nodes and inspect the running queries:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 |
-[ RECORD 7 ]----+----------------------------------------------- datid | 17273 datname | test pid | 22237 leader_pid | usesysid | 10 usename | hs application_name | citus_internal gpid=10000022078 client_addr | 127.0.0.1 client_hostname | client_port | 37460 backend_start | 2025-01-08 21:20:52.062818+01 xact_start | 2025-01-08 21:20:52.066472+01 query_start | 2025-01-08 21:20:52.807719+01 state_change | 2025-01-08 21:20:52.807722+01 wait_event_type | wait_event | state | active backend_xid | backend_xmin | 2491 query_id | 6772522768391086066 query | SELECT key_id, count(*) AS count FROM public.t_data_106302 t_data WHERE true GROUP BY key_id backend_type | client backend |
Citus sets a proper application_name, so standard PostgreSQL tooling is enough to observe cluster activity.
The limits of brute force
Not every query survives this scale. Consider a computation that, for each key_id, derives the average distance between ascending values:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 |
test=# explain SELECT key_id, avg(data - lag) FROM ( SELECT *, lag(data, 1) OVER (PARTITION BY key_id ORDER BY data) FROM t_data ) AS x GROUP BY 1; QUERY PLAN ----------------------------------------------------------------------------------- Custom Scan (Citus Adaptive) (cost=0.00..0.00 rows=100000 width=36) Task Count: 32 Tasks Shown: One of 32 -> Task Node: host=localhost port=5432 dbname=test -> GroupAggregate (cost=6782453812.10..7675036577.03 rows=317 width=36) Group Key: t_data.key_id -> WindowAgg (cost=6782453812.10..7431604910.98 rows=32457554944 width=20) -> Sort (cost=6782453812.10..6863597699.46 rows=32457554944 width=12) Sort Key: t_data.key_id, t_data.data -> Custom Scan (ColumnarScan) on t_data_106299 t_data (cost=0.00..6326942.55 rows=32457554944 width=12) Columnar Projected Columns: key_id, data (12 rows) |
The whole dataset must be sorted, passed through a window function, and then grouped — for 1 trillion rows, on a local machine. Scale also applies to expectations: once datasets grow, some operations become markedly harder and cannot be run carelessly.



