Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

iTechGuides is reader-supported. When you buy through links on our site, we may earn an affiliate commission. As an Amazon Associate I earn from qualifying purchases. Learn more

A single PostgreSQL server can be outgrown, and the way out does not have to mean leaving PostgreSQL. The right step depends on which resource is actually saturated: query plans, CPU, memory, storage, table size and retention, read traffic, availability requirements, or sustained write throughput. Each of those points to a different remedy, and the remedies are not interchangeable. Partitioning, replication, logical replication, and distributed PostgreSQL such as Citus solve different problems, so the first job is to diagnose the bottleneck.

Start by identifying the actual bottleneck

“We have outgrown Postgres” usually describes one of several distinct conditions. A slow report that scans a 400 GB table, a connection pool that is saturated at peak, a primary that cannot fail over quickly, and a write rate that the storage subsystem cannot absorb all look similar from the application side, but they need different fixes. Treating them as one problem tends to produce an expensive architecture change that leaves the original slow query in place.

Work through the following checks before changing topology:

Free tools Windows power users keep installed

One-click scans. No signup required.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
  • Query plans. Run EXPLAIN (ANALYZE, BUFFERS) on the slowest statements. Look for sequential scans over large tables, row estimates that are far from actual counts, and sorts or hashes that spill to disk. Many “database too small” complaints end here.
  • Resource saturation. Compare CPU, memory, disk IOPS, and disk throughput against their limits during the slow periods. A workload that is fully I/O-bound behaves very differently from one that is CPU-bound.
  • Connection count. Check whether latency is driven by the number of active backends rather than by any single query. Connection pooling is often the cheapest relief.
  • Read or write mix. Determine whether the pressure comes from reads that could be served elsewhere or from writes that must land on one primary.
  • Availability requirement. Establish how long an outage may last and how much data loss is tolerable. These are business requirements, and they determine whether replication is needed at all.
  • Retention and table lifecycle. Decide whether old rows are deleted in bulk, archived, or kept indefinitely. Bulk deletes against huge tables are a common hidden cost.

Only after these questions have answers should you choose among the options below. The comparison table later in this article maps each constraint to a remedy.

Tune queries and use parallel query carefully

Query and index work is the cheapest intervention and it does not change deployment topology. Adding a targeted index, rewriting a correlated subquery, or removing an unnecessary join can produce gains that no hardware change would match. These gains are specific to the queries you change, so verify them against production-like data volumes.

PostgreSQL’s parallel query can also speed some eligible reads, but it is not a general scaling switch:

  • The planner does not generate parallel plans for statements involving writes or row locking.
  • Operations classified as parallel-unsafe disable parallel query for that statement.
  • Each parallel worker is a separate process. The PostgreSQL resource consumption documentation notes that a query using four workers can use up to five times the resources of the same query run without workers. Under concurrent load, that multiplication can slow the whole server down.

Treat the worker setting as a concurrency parameter to tune against measured results, not as a way to add capacity.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Native partitioning: one table, many physical pieces, one server

Declarative partitioning splits one logical table into several physical tables inside the same PostgreSQL database system. In the PostgreSQL 18 documentation, the partitioned parent is a virtual table with no storage of its own; each partition is an ordinary table with defined bounds, and inserts are routed to the matching partition automatically.

Partitioning helps in two situations:

  • Partition pruning. When a query filters on the partition key, the planner can skip partitions that cannot contain matching rows. A time-partitioned events table queried for the last seven days touches only a few partitions.
  • Lifecycle operations. Dropping or detaching an old partition is typically far cheaper than deleting millions of rows from a monolithic table, and it avoids the bloat that large deletes leave behind.

It does not distribute writes across machines. Every partition still lives on the same server and draws on the same CPU, memory, and disk. The documentation also warns that planning overhead and memory use grow when many partitions remain relevant to a query, and that more partitions are not automatically better. Choose a partition key that matches your dominant query filters, and keep the partition count moderate.

Replicas: availability and read distribution

PostgreSQL’s high-availability documentation describes servers that cooperate so that a standby can take over when the primary fails, or so that several machines can serve the same data. It is explicit that different solutions handle synchronization differently and that no single approach removes every trade-off.

In practice, a physical standby helps in two ways:

  • Failover. A standby that is kept current can be promoted when the primary fails. Your recovery time depends on how quickly failure is detected and how the promotion is automated.
  • Read offload. Read-only queries that tolerate slightly stale data can run on standbys, which reduces load on the primary.

Replicas do not add write capacity. All writes still go to the primary. Read-only traffic sent to a replica can also see replication lag, so applications that need read-your-own-writes behavior must route those reads to the primary or wait for the replica to catch up. Synchronous versus asynchronous replication is the central decision: synchronous modes protect data at the cost of commit latency, while asynchronous modes keep the primary fast but can lose the most recent transactions on failover.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Logical replication: copying selected data, not sharding writes

Logical replication works at the level of changes to tables, using publications on the source and subscriptions on the target. A typical subscription first copies a snapshot of the existing table data, then continuously sends subsequent changes. Within a single subscription, changes are applied in the order the publisher committed them.

The PostgreSQL documentation lists several uses:

  • Replicating a subset of tables or rows rather than the whole cluster.
  • Consolidating data from several databases into one analytical store.
  • Replicating between major PostgreSQL versions during an upgrade.
  • Sharing data between databases.

It has real prerequisites. The publisher must run with the logical wal_level, replication slots must be managed so that unconsumed WAL does not accumulate, and subscriptions consume background worker capacity. Logical replication is a data-movement mechanism. It is not a multi-writer cluster, and it will not transparently route an application’s writes across several nodes.

Distributed PostgreSQL: Citus and multi-node tables

Citus is an extension that turns a group of PostgreSQL nodes into a distributed database. Its documentation describes distributed tables sharded across the nodes of a cluster, reference tables replicated to every node so that joins against them stay local, and a distributed query engine that routes or parallelizes statements across the cluster. Microsoft’s Citus documentation describes the same architecture for its managed offering.

This is the option that actually adds write and storage capacity beyond one machine, but it is only a good fit when the data and the queries cooperate:

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
  • A clear distribution column exists that most large tables share, and most queries filter or join on it.
  • Cross-node joins, transactions, and unique constraints that do not include the distribution column are rare or acceptable at the application level.
  • The team can operate a multi-node system, including rebalancing shards and handling node failures.

Adopting a distributed layer usually requires schema changes, such as choosing distribution columns and adjusting constraints, and it may require application changes for queries that the planner cannot route efficiently. Verify the specific Citus version and hosting service you intend to use, because feature sets and limits differ across releases and platforms. No benchmark in the sources supports a general claim of speedup for any given workload, so measure your own schema and queries before committing.

Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

How the options compare

Measured constraint What to investigate first What it changes Main trade-off
Inefficient plans or a few expensive reads Query plans, indexes, schema and query rewrites, eligible parallel query Nothing in the deployment topology Gains are query-specific; parallel workers add CPU, memory, and I/O load
Large table with time- or key-bounded access and retention Declarative partitioning Table design within one server Poor partition keys or too many relevant partitions hurt planning and memory use; no added write capacity
Availability or more read capacity Physical standbys, read routing, failover automation Adds servers holding the same data Synchronous versus asynchronous behavior, replication lag, and failover handling; writes stay on the primary
A subset of data, or a downstream analytical copy Logical replication Adds a subscriber that receives selected changes Requires logical WAL level, replication slot management, and worker capacity; not a general write-sharding layer
Write or storage capacity beyond one node, with distributable data and queries Distributed PostgreSQL such as Citus Distributes tables and queries across nodes Distribution keys, cross-node operations, and operational complexity; verify version and service limits
Operations burden rather than an engine limit A managed PostgreSQL service Shifts routine operations to the provider Feature sets, limits, and pricing vary by provider; not stated in this article

When comparing real options, evaluate five things: which bottleneck the option addresses, whether it forces changes to application or schema assumptions, its consistency and failover behavior, its operational complexity, and its compatibility with the PostgreSQL features and extensions your application already depends on.

What the official limits do and do not tell you

The PostgreSQL 18 documentation lists a database size with no hard limit, and a maximum relation size of 32 TB when the default 8 KB block size is used. Those are hard ceilings, not operating targets. The same documentation warns that practical limits from performance or available disk space can apply much earlier.

There is no universal row count or request rate at which a team must leave a single node. A table with a billion rows can be fast on well-indexed queries and slow on unindexed scans, and a smaller table with heavy write contention can be the bigger problem. Use a representative workload, realistic data volume, and your actual latency and availability targets to decide.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

A practical order of operations

  1. Capture the slowest queries and the resource metrics from the slow periods.
  2. Fix plans, indexes, and connection handling, then measure again.
  3. If a large table is the problem and access or retention is time- or key-bounded, partition it.
  4. If availability or read load is the problem, add standbys, define the synchronization mode, and test failover.
  5. If you need a subset or a copy for another system, use logical replication and confirm the slot and worker configuration.
  6. Only when writes or storage must exceed one machine and the data and queries can be distributed should you evaluate Citus or a comparable distributed PostgreSQL system.

Following this order keeps you on PostgreSQL throughout, and it means that each architectural change is justified by a bottleneck you have already measured.

“

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.