Find expensive work rather than blaming table size
A hundred-million-row table may answer an indexed lookup quickly. A smaller table can stall because of a poor plan, long transaction or contention on one record. Data volume matters, but does not by itself choose the next intervention.
Collect query frequency and duration, read volume, connection waits, locks and disk use. A 200 ms query repeated ten thousand times can consume more than an occasional five-second report. Break down the workload by journey and customer.
For PostgreSQL, EXPLAIN ANALYZE compares planner estimates with actual rows and timings. Large differences help expose inaccurate assumptions about data. ANALYZE executes the query, so inspecting modifying statements requires deliberate care and an appropriate environment.
We begin with the period users experienced the problem and compare ordinary operation. Did frequency, result size, active transactions or waits change? A daily average can hide a ten-minute peak. Choosing the measurement window is part of diagnosis, not just a dashboard display preference.
We rank queries by total time and individual latency separately. One identifies overall resource consumption, the other slow user paths. Frequent short queries can be valuable targets even when no individual customer notices them. The two rankings answer different questions and need not point to the same SQL.
Plans are checked for actual rows, repeated passes and intermediate result size. A JOIN that behaves well on a small selection can become expensive for a large account. We try representative parameter values, including popular filters and skewed customers, rather than use a convenient single ID.
Lock waits need the owner and the reason it holds the resource. The waiting query looks expensive in elapsed time while doing little execution. An index on it will not release another transaction's lock. We investigate why the holder stays active and whether its critical section can be shorter.
We decide how to measure safely. Even a SELECT can impose substantial load, and modifying statements affect records. A suitable copy or bounded experiment comes first, with differences in data and resources recorded. Diagnosis should establish a cause without recreating an uncontrolled outage for customers.
Optimize database queries before scaling
An order page loads 50 rows, then fetches each customer separately: 51 database calls, a familiar N+1. Batched loading, an appropriate JOIN and a bounded result set can help more than a larger server. The improvement also keeps paying off as traffic grows.
Choose indexes for actual filters, ordering and data distribution. An index need not serve every query; scanning a large portion sequentially can be cheaper. Check read performance, write overhead, index size and maintenance after the change.
A transaction waiting for an external service holds resources during someone else's work. Revisit its boundary, then explicitly preserve workflow integrity outside it. Simply making a transaction shorter does not solve the new partial-failure cases.
Fixing N+1 should not produce one enormous JOIN. Hundreds of items per order can multiply rows and transfer size. Sometimes two batched queries are better, sometimes a bounded JOIN is. We compare actual work and response volume instead of treating SQL statement count as the only objective.
A composite index serves particular predicates and ordering. It cannot help every filter combination equally. If the UI permits dozens of arbitrary sort paths, we prioritize important ones or constrain the product. Indexing every combination becomes expensive for writes and maintenance long before it solves every query.
Pagination affects growth. Deep offsets can require progressively more work skipping earlier rows. A cursor over a stable ordering may fit sequential browsing better. It must handle ties, newly arriving records and user expectations. Changing pagination semantics is a product decision as well as a query optimization.
Transaction length follows the business path. Holding a database lock while a customer enters card details is a poor boundary. A short reservation with an expiry and an external payment stage can work, but needs explicit confirmation and expiration rules. Removing the lock does not remove the need for correctness.
We inspect maintenance too: statistics, space reclamation, index growth and long schema changes. Adding indexes or columns can compete with live traffic. The method follows the chosen database's capabilities and is tested for impact. A faster final query does not justify hours of blocking during the migration that produced it.
A cache requires a freshness contract
Delivery reference data is easier to cache than the last available item. Define acceptable staleness, update behavior and misses for each object. If a displayed price is old, decide which price checkout will confirm.
Estimate database load without the cache. If it serves 90% of reads, losing it can increase database read traffic roughly tenfold at unchanged arrivals. This is a calculation for the example, not a universal cache property. Bound refresh concurrency, spread expiration and decide when an old response is acceptable.
For price data, display and order confirmation are separate decisions. A page may show a value within an agreed freshness window, while checkout verifies current terms. If the price changed, the customer sees the new total before committing. This makes caching part of the workflow rather than merely a TTL setting.
Invalidation follows ownership. If several systems can change a price, deleting a key in one service may miss another source. Versions or an agreed change stream can help, but the source and missed-event behavior need rules. Expiry bounds age; it does not promise immediate freshness.
A popular key expiring can send hundreds of requests to refresh the same value. Limiting refreshers and temporarily serving an acceptable stale response can reduce the stampede. Refresh failure still needs visibility, or the cache can conceal a broken dependency while appearing to keep the service available.
We inspect hit rate by important data class. An overall ninety-five percent can come from cheap reference data while the expensive path misses every time. Value is measured in database work removed and user latency, not merely the proportion of keys found in memory.
Caching need not come first. It adds another store and states such as missing, stale, refreshing and unavailable. If a good index and a small response already meet the budget, invalidation complexity may cost more than it saves. We choose it after measuring a specific expensive path and acceptable freshness.
A read replica changes when writes become visible
An asynchronous replica receives changes later. A customer saving an address and reading it immediately can see the previous value on a lagging replica. That looks like a lost write even when the primary saved it. Route this path to the primary or establish another explicit visibility condition.
Synchronous replication has distinct acknowledgement modes. Receiving a record and applying it are different guarantees. Stricter confirmation can add latency and dependence on replica availability. Choose it for required operations rather than treating consistency as one setting.
Read replicas offload reads; they do not split the write workload. We watch replication lag, replay capacity and failover behavior. A replica is not a backup either: an accidental deletion can reach every copy quickly, so recovery needs its own test.
Read routing is defined per operation. Catalogs and older reports can tolerate lag; immediate confirmation of an address change usually cannot. Reading the primary for a few sensitive paths may be simpler than a universal replica-wait mechanism. Those reads remain in the primary's capacity budget.
A replica is not free to the primary: changes are logged, sent and applied elsewhere. Heavy write load or slow replay requires inspecting the whole path rather than expecting linear capacity from every additional copy. We also test catch-up after substantial lag, when normal steady-state behavior no longer describes it.
Failover needs ownership rules: who selects the new primary and how the previous one is prevented from accepting writes. Two nodes both believing they own writes can diverge. Changing a connection address does not solve that. The rules must match the actual cluster-management mechanism.
The application participates in failover. Connections break, some results become unknown and pools reconnect. We check timeouts and safe repeats that preserve the business operation. A wave of simultaneous reconnections can overload a weakened cluster, so client recovery behavior belongs in the test.
Backup recovery checks data completeness, restoration point and external reconciliation as well as database startup. Another system may have confirmed operations independently. Replicas, backups and reconciliation protect different parts of the outcome even when all are loosely described as data protection.
Table partitioning and sharding have different costs
Partitioning divides a logical table into pieces. Monthly event partitions can simplify removing old periods and reduce the data visited by a weekly query. Queries must permit pruning for that benefit. Dividing a table does not imply placing it on independent servers.
Sharding distributes data among databases. A company key can keep common operations local, but one large company may overload its shard. Another key may balance load while turning ordinary reports into cross-shard queries.
Before moving data, resolve identifiers, uniqueness, cross-shard transactions, queries without the key, rebalancing and recovery. If these operations dominate the product, distribution may cost more than the expected capacity gain.
Time partitions fit an event lifecycle when queries usually know the relevant period. A query without a time predicate may still traverse many parts. We collect actual filters before choosing boundaries. Partitioning for its own sake can increase maintenance without making important requests meaningfully faster.
Before sharding, we estimate activity as well as storage. Ten similarly sized parts can differ enormously in traffic. A large customer, popular category or bulk operation breaks a simple balance. We need visibility and a redistribution path, not just a count of equally sized databases.
The shard key is checked against frequent operations. Keeping orders and items together simplifies local transactions. Splitting payment and order creates a cross-database protocol for ordinary confirmation. The decision follows required guarantees and the amount of extra application logic, rather than a generic desire to distribute data.
A query without a shard key may need every part. More parts increase cost and the chance one delays the combined result. Search or analytics can need another data path, whose freshness and recovery must also be defined. Moving a query elsewhere does not make its guarantees disappear.
We test moving a customer between shards before large-scale growth. Data must be copied, concurrent changes retained and routing switched consistently. If that requires a long pause, the condition should be known. Additional storage capacity is useful only alongside a workable way to manage its distribution.
A change must survive migration and operations
Compare the representative workload before and after: latency, throughput, contention, cost and failure headroom. Changing queries, hardware and placement simultaneously without attribution makes the result hard to explain.
The migration needs data completeness checks, write cutover rules and a response to mismatches. Rollback must account for writes that arrived after switching; restoring the old schema may no longer be enough. We are finished when the team can operate and recover the new setup, not when the copy command succeeds.
We choose completeness checks before moving data. Equal row counts do not mean equal content. IDs, important totals, relationships and selected records need comparison. Business invariants guide checks so a successful copy command cannot hide a migration that changed the meaning of the records.
When writes continue, changes after the initial copy must arrive too. We define ordering, duplicate handling and lag detection. Writing to two stores is not safe just because both calls exist: one can succeed while the other fails. Divergence needs detection and a correction path.
Cutover is a sequence with owners: verify catch-up, restrict writes if needed, switch readers, check the product and continue traffic. Each step has a stop condition. "The script finished" describes technical completion, not data quality or whether customers can correctly use the new path.
Rollback includes new writes. Once customers use the new schema, returning to the old path requires compatibility or transferring their changes. Sometimes disabling a feature briefly and fixing forward is safer. That choice is made before the incident, when the team has time to evaluate it.
After cutover, we observe tail latency, errors, growth and maintenance. The old path is removed after checks rather than at the same instant as switching. A completed migration means correct data, a working product and a team that can operate and recover the new layout.
Choosing the next step for order data
The order-list page slows as the database grows. We compare periods and separate connection waits from execution. N+1 and unnecessary fields appear. Fixes are tested across small and large customers, then writes and maintenance are checked so a new index does not consume capacity needed elsewhere.
If stable repeated reads remain costly, we evaluate caching with freshness rules and cache-loss behavior. Immediate reads after updates keep a separate path. Reports that tolerate lag can use a replica, with replay and failure tests. Another copy is not presented as protection from every data-loss scenario.
Partitions follow actual time filters and history deletion. Shards follow workload distribution and cross-shard costs. Each component gets write ownership, recovery and migration rules. Added capacity must remain manageable by the actual team. Otherwise storage scale hides an operational limit that will emerge at the next incident.
Cutover is ready when checks agree, new writes remain intact, the product works and mismatch handling is executable. We observe traffic before removing the old path. The next step may again be local SQL: every improvement moves the constraint. That is a normal process of scaling, not a reason to jump immediately to the most distributed design.