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.