TL;DR
- Migrate from PostgreSQL only when you have already exhausted optimisation, indices, partitioning, read replicas and workload separation.
- The right reason for a distributed database is usually one of these: multi-region writes, regional availability, data volume per shard, or real throughput limits.
- CockroachDB, YugabyteDB, Google Spanner, Citus and Vitess solve different problems. They are not “bigger Postgres”.
- The main cost is not the licence. It is latency, consistency, debugging, migrations, observability and a change in the team’s mental model.
- If your problem is a slow query or missing indices, a distributed database will make the problem more expensive, not simpler.
The common mistake: confusing growth with the need for distribution
PostgreSQL can handle more than most teams imagine. A well configured PostgreSQL 16, with the right indices, tuned autovacuum, reviewed queries and decent hardware, can serve many B2B SaaS products, internal platforms and eCommerce applications without needing a distributed database.
The right question is not “when does Postgres stop scaling?”. The right question is: “what exactly is the limit we are hitting?”.
There are very different limits:
- CPU saturated by poorly written queries.
- I/O saturated by missing indices or excessive writes.
- Locks caused by long transactions.
- Large tables without partitioning.
- High latency caused by clients in another region.
- Need for active writes in multiple regions.
- Maintenance window too short for DDL operations.
- Lack of isolation between analytical and transactional workloads.
Only some of these justify a distributed database.
In practice, I see too many teams jump to CockroachDB, YugabyteDB or Spanner because “we are going to grow”. That is rarely a good criterion. Distribution should be a response to a concrete, measured and repeatable constraint. If you cannot show charts for CPU, IOPS, locks, p95, p99, table size, WAL volume and incident frequency, you do not yet have a diagnosis. You have architectural anxiety.
Before migrating, exhaust normal Postgres
Before thinking about a distributed database, there is a list of serious PostgreSQL work that usually solves 80% of cases.
First, indices. Not “put indices on everything”, but indices aligned with real queries. pg_stat_statements is mandatory. Without it, you are optimising by guesswork. In PostgreSQL 16, you should look at queries by total time, average time, calls and disk reads.
SELECT
query,
calls,
round(total_exec_time::numeric, 2) AS total_ms,
round(mean_exec_time::numeric, 2) AS mean_ms,
rows
FROM pg_stat_statements
ORDER BY total_exec_time DESC
LIMIT 20;
Second, partitioning. A table with 500 GB is not automatically a problem, but a table with 500 GB, frequent deletes, bloated indices and time range queries can be an avoidable problem. Partitioning by month or by tenant can reduce scans, make retention easier and make maintenance more predictable.
Third, read replicas. Many systems mix transactional traffic with dashboards, CSV exports, reports and external integrations. Moving heavy reads to replicas can reduce pressure on the primary without changing the data model. Pay attention to replication latency. If the application assumes immediate reads after writes, you can introduce subtle bugs.
Fourth, pooling. In SaaS with Node.js, serverless or many workers, I have seen databases suffer more from too many connections than from complex queries. PgBouncer in transaction pooling can be the difference between 300 chaotic connections and a stable load. But it also has trade-offs, for example prepared statements and stateful sessions.
Fifth, hardware and configuration. It sounds unsophisticated, but moving from a small instance to a machine with NVMe, enough RAM for the working set and predictable IOPS is often cheaper than rewriting the data layer. In cloud, the difference between general purpose storage and provisioned storage can be brutal at p99.
My rule: if you do not yet have pg_stat_statements, tested backups, autovacuum metrics, replication lag alerts and explanations for the 20 most expensive queries, you are not ready to choose a distributed database.
Real signs that Postgres may be at the limit
There are signs that justify seriously opening the discussion.
The first is writes that do not fit on a single primary. Traditional PostgreSQL scales reads relatively well with replicas, but writes remain concentrated on the primary. If you have thousands of writes per second, contention on indices, very high WAL and you cannot reduce load through batching or redesign, distribution may make sense. Even so, it is worth separating “high writes” from “poorly modelled writes”. Inserting append-only events is different from updating the same rows concurrently.
The second is active multi-region. If you have users in Europe and the United States who need to write locally and require p99 below 150 ms, a central database in eu-west-1 is not enough. Physical latency rules. Lisbon to Virginia can easily exceed 80 ms just on the network. With TLS, multiple queries and transactions, p99 rises quickly. Here, a database such as Google Spanner or CockroachDB can reduce local latency, but it will force you to think about consistency, conflicts and data placement.
The third is regional availability. If your requirement is to keep accepting writes when a cloud region fails, PostgreSQL with a synchronous replica in another region can work, but the latency cost is high and failover is operationally delicate. Distributed databases were designed for this type of scenario, although they do not eliminate difficult decisions.
The fourth is operational size. There is no magic number. 1 TB in PostgreSQL can be healthy. 200 GB can be chaos. But when backups, restores, vacuum, reindexing, migrations and schema changes no longer fit within operational windows, you need to reassess the architecture. Sometimes the answer is partitioning or archiving. Other times it is sharding or distribution.
The fifth is tenant isolation. In B2B SaaS, the time may come when a few large customers disrupt everyone else. You can solve this with partitioning, separate databases per tenant, separate schemas, or application-level sharding. A distributed database can help, but it is not always the first option.
Concrete alternatives and when they make sense
CockroachDB is attractive for teams that want SQL, distributed transactions and partial compatibility with PostgreSQL. It is a good candidate when you need multi-region, availability and automatic data distribution. But it is not PostgreSQL. There are differences in extensions, query behaviour, optimisations and transaction latencies. The official documentation at https://www.cockroachlabs.com/docs/ should be read before assuming compatibility.
YugabyteDB follows a similar line, with a PostgreSQL-compatible API through YSQL and a distributed architecture based on tablets. It can be interesting for distributed SQL workloads, especially when compatibility with Postgres is important. But, as with CockroachDB, distributed transactions have a cost. Consult https://docs.yugabyte.com/ before deciding.
Google Spanner is probably the most mature option for global consistency with distributed SQL, especially on Google Cloud. TrueTime is a real technical advantage. But you are buying a platform, not just a database. The cost model, primary key design and cloud dependency should be evaluated coldly. Documentation: https://cloud.google.com/spanner/docs.
Citus, now part of the Microsoft ecosystem, is an extension that distributes PostgreSQL across shards. It makes sense when you want to keep much of the Postgres model and have a clear distribution key, for example tenant_id. It is excellent for some multi-tenant SaaS and distributed analytical workloads. It is less elegant when queries frequently cross tenants or when there is no good distribution key. Documentation: https://docs.citusdata.com/.
Vitess is another category. It was created to scale MySQL, not PostgreSQL. It is used in large-scale environments and supports sharding, routing and online operations. It makes sense if you are in the MySQL world and have a team to operate that complexity. It is not a direct answer for anyone using PostgreSQL, but it is important as an architectural reference. Documentation: https://vitess.io/docs/.
The practical comparison:
| Option | Best case | Hidden cost |
|---|---|---|
| PostgreSQL 16 with tuning | Normal B2B SaaS, moderate eCommerce, back office, internal fintech | Single primary for writes |
| PostgreSQL with partitioning and replicas | Heavy reads, time-based tables, separate reporting | Application complexity and replica lag |
| Citus | Multi-tenant with a good distribution key | Cross-shard queries and data planning |
| CockroachDB or YugabyteDB | Distributed SQL, availability, multi-region | Transaction latency and differences from Postgres |
| Spanner | Global consistency on Google Cloud | Cost, schema design and platform dependency |
What changes technically in a distributed database
The most important change is that the network becomes part of the database. In classic PostgreSQL, a transaction is local to the server. In a distributed database, a transaction may touch several nodes, regions and internal consensus mechanisms. This affects p99.
A query that takes 15 ms in PostgreSQL may take 80 ms if it touches multiple shards or ranges. A transaction with five statements may stop being cheap if it involves distributed coordination. The problem is not the average. It is p95 and p99 under real load.
Key design also changes. Sequential IDs can create hotspots. Random keys can distribute writes well, but harm locality. Keys composed of tenant_id and time can be excellent for SaaS, but poor for global queries. There is no free lunch.
Consistency is another change. Many distributed databases promise ACID transactions, but you need to understand the scope. Transactions within the same shard can be cheap. Cross-shard transactions can be expensive. “Stale” reads may be acceptable for dashboards, but not for payments, balances, critical inventory or permissions.
In fintech and payments, for example, I would be conservative. I would not distribute financial writes without a very strong reason, failure testing and independent reconciliation. For balances, ledgers and charges, operational simplicity is worth a lot. Stripe, for example, publishes extensive documentation about idempotency and webhooks at https://docs.stripe.com/, and that discipline is as important as the database chosen.
You also have to rethink migrations. ALTER TABLE that used to be trivial can become a long operation. Global indices may have limitations. Backfills can saturate clusters. And observability stops being “look at the database CPU”. You need metrics per node, shard, range, region, internal queues, retries, conflicts and consensus latency.
Production pitfalls that cost dearly
The most common pitfall is using a distributed database as if it were PostgreSQL with more machines. The application keeps large joins, broad transactions, queries without a distribution key and back-office jobs that scan entire tables. Result: more nodes, more latency, more cost and less predictability.
Another pitfall is ignoring retries. In distributed systems, transient errors are normal. Transactions can fail because of contention, a leaseholder change, a serialisation conflict or a timeout. The application must know how to retry idempotent operations. If your code does not distinguish a permanent error from a transient error, you will have strange incidents.
Short example of an essential pattern in TypeScript pseudocode:
for (let attempt = 1; attempt <= 3; attempt++) {
try {
return await runTransaction()
} catch (err) {
if (!isRetryableDatabaseError(err) || attempt === 3) throw err
await sleep(50 * attempt)
}
}
The important detail is not the snippet. It is the discipline: repeatable operations, idempotency keys, explicit timeouts and logs with correlation id.
Another gotcha: timestamps. In multi-region systems, blindly relying on now() to order business events can be dangerous. Some databases provide strong guarantees, others do not. Even with guarantees, the semantics of “happened before” must be modelled. For critical events, use versions, per-entity sequences or append-only ledgers.
There is also the cost pitfall. A managed PostgreSQL can cost a few hundred euros per month for a growing product. A multi-region distributed cluster, with traffic between regions, observability and overprovisioning for failures, can multiply that several times. The cost is acceptable when it buys availability or latency that the business really needs. It is waste when it buys psychological comfort.
A decision process that works
I would use a simple sequence.
First, measure. Collect 14 to 30 days of metrics: CPU, memory, IOPS, table and index size, WAL generated per hour, most expensive queries, locks, deadlocks, replication lag, p95 and p99 per critical endpoint. Without this, the decision is weak.
Second, classify the problem. Is it read, write, geographical latency, availability, data size, customer isolation, or operation? Each category points to different solutions.
Third, try the least distributed solution. Indices, query rewrite, partitioning, replicas, queues, cache, cold archive, materialized views, read models, or OLTP/OLAP separation. Often, moving reporting to ClickHouse, BigQuery or DuckDB over object storage solves more than changing the transactional database.
Fourth, run a test with a real workload. Do not use only synthetic benchmarks. Reproduce the 20 most important queries, background jobs, write spikes and migrations. Measure p50, p95, p99 and retry rate. A good initial target might be p99 below 200 ms for critical operations within the same region, but it depends on the product.
Fifth, design failures. Turn off a node. Simulate loss of a region. Increase latency between regions. Deploy during a backfill. Restore a backup. If the team cannot operate these scenarios in a test environment, it is not ready for production.
Sixth, calculate total cost. Include licences, cloud, inter-region traffic, observability, engineering time, training, migration, temporary dual writes, rollback and support. The question is not “which technology is better?”. It is “what is the smallest system that meets the requirements for 18 to 24 months?”.
My opinion is simple: PostgreSQL should be the default choice until proven otherwise. A distributed database comes in when there is a constraint that Postgres cannot solve without creating more risk than it removes. Not before.
Conclusion
Migrating from PostgreSQL to a distributed database is an architectural decision, not a rite of passage. If the problem is local, solve it locally. If the problem is multi-region, strong availability or writes beyond what a primary can support, then it is worth evaluating alternatives.
The worst migration is the one that replaces a known database with a poorly understood distributed system. If you are facing a similar problem, book a call at https://impact-origin.com/agendamento.