Database choices are expensive to undo. Data outlives the code around it, every service grows assumptions about its store's consistency and query model, and migrating a live system is a project measured in quarters. Even so, most choices are made by familiarity, by a conference talk, or by a benchmark that looks nothing like the real workload.

This guide gives a repeatable method. You describe the workload as numbers first, then let those numbers eliminate options. It covers the questions that decide the choice, the main families and what each charges you for, a worked sizing example with code, how to benchmark, when a second store is justified, and the failure modes behind most painful migrations. For deep dives on specific engines, see When to pick PostgreSQL, When to pick Cassandra and When to pick DynamoDB.

Advertisement

Start from the workload, not the product

A database is a set of trade-offs packaged together: a data model, a query model, a consistency model, a storage engine and an operational model. Products differ in which trade-offs they make, so the only way to compare them is against a description of what you need. Write that description before anyone names a product.

A sensible default exists. For a new system whose requirements are still moving, a mainstream relational database is the lowest-risk start: it supports ad hoc queries you have not thought of yet, enforces invariants with transactions and constraints, and a single well-sized node with replicas serves far more load than most products ever see. The method below is mostly a disciplined way of asking whether any measured requirement rules that default out.

The questions that decide it

  1. What invariants span more than one row? Balances that must never go negative, stock that must not be oversold, usernames that must be unique. These need transactions or conditional writes in the store. Pushing them into application code over an eventually consistent store is where data corruption comes from.
  2. How is data accessed? List every query: point lookups by key, range reads, ad hoc filters, joins, aggregations over large ranges, full-text search, similarity search. Stores built for key access are fast and scalable precisely because they refuse arbitrary queries.
  3. How big and how fast? Data size after growth, peak writes per second, peak reads per second, and how these split by key. Hot keys matter as much as totals.
  4. What latency, availability and geography? A p99 objective, a recovery time and acceptable data loss on failure, and whether writes must be accepted in more than one region at once.
  5. Who will run it? Backups, restore drills, upgrades, schema changes, monitoring and on-call. A managed service shifts much of this; a self-run cluster needs people who know it.
A workload-first decision flowQuery inventoryevery read and write, with ratesMulti-row invariants?money, stock, uniquenessyes or unsureStart relationalone node + replicasnoKey-access only?known queriesKey-value / wide-columnmodel per querySized past one node?arithmetic, then benchmarkwritesDistributed storesharded SQL or NoSQLno: keep itAdd specialised stores only for measured needsanalytics, search, time series, vectors: fed by CDCEach arrow is a question you answer with numbers from the inventory, not with a product name.
Questions in order. Most workloads stop at the relational box; the others need a measured reason.
Advertisement

Families and what each charges you

FamilyGood atWhat you payReach for it when
Relational (PostgreSQL, MySQL)Transactions, constraints, joins, ad hoc SQLWrite scaling beyond one primary needs sharding or a distributed SQL productDefault for system-of-record data
Distributed SQLSQL and transactions across nodes and regionsHigher per-query latency, cost and operational complexityRelational needs that measurably exceed one primary
Key-value and wide-column (DynamoDB, Cassandra)Predictable latency at very high write ratesNo joins; you model one table per query and denormaliseKnown, stable, key-based access at scale
Document (MongoDB and others)Nested, varied records read as a unitCross-document invariants and reporting are harderAggregates naturally read and written whole
Columnar analyticalScans and aggregations over billions of rowsPoor at single-row updates and high-rate point readsReporting, BI, product analytics
Time-seriesAppend-heavy metrics, retention, downsamplingNarrow query modelTelemetry, monitoring, IoT
Search engineFull-text relevance, facetingNot a system of record; eventual indexingUser-facing search over text
Vector indexApproximate nearest-neighbour searchApproximate results; recall versus latency tuningSemantic search and retrieval

The split between transactional and analytical work is the most common reason to have two stores; OLTP versus OLAP explains why one engine rarely does both well at scale.

Worked example: an order service

Suppose you are building an order service. Write the query inventory as code, so it can be reviewed and rerun when forecasts change:

# query_inventory.py: turn the access patterns into numbers before naming a product
from dataclasses import dataclass

@dataclass
class Q:
    name: str
    kind: str          # "point", "range", "scan", "write"
    peak_per_s: float
    p99_ms: float      # latency objective
    key: str           # what the query is looked up by

INVENTORY = [
    Q("create order",          "write",  300,  50, "order_id"),
    Q("get order",             "point", 2000,  20, "order_id"),
    Q("orders for customer",   "range",  900,  50, "customer_id, created_at"),
    Q("update status",         "write",  450,  50, "order_id"),
    Q("daily revenue by region", "scan", 0.01, 60000, "created_at"),
]

ROW_BYTES, INDEX_FACTOR = 2_000, 2.5       # row size incl. overhead; indexes and bloat
ROWS_PER_YEAR, YEARS, GROWTH = 20e6, 5, 1.5  # plan for 1.5x the forecast

storage_gb = ROWS_PER_YEAR * YEARS * GROWTH * ROW_BYTES * INDEX_FACTOR / 1e9
writes = sum(q.peak_per_s for q in INVENTORY if q.kind == "write")
reads = sum(q.peak_per_s for q in INVENTORY if q.kind in ("point", "range"))
scans = [q.name for q in INVENTORY if q.kind == "scan"]

print(f"storage after {YEARS} years: {storage_gb:,.0f} GB")
print(f"peak writes/s: {writes * GROWTH:,.0f}  peak key reads/s: {reads * GROWTH:,.0f}")
print("scan workloads (consider a separate analytical copy):", scans)

Running it prints:

storage after 5 years: 750 GB
peak writes/s: 1,125  peak key reads/s: 4,350
scan workloads (consider a separate analytical copy): ['daily revenue by region']

Now read the numbers. 750 GB after five years with growth margin fits comfortably on one modern database node. About 1,100 writes per second is well within what a single relational primary sustains for small transactional rows, though you confirm that by benchmarking. Every hot query has a clear key, so indexes on order_id and (customer_id, created_at) serve them. Creating an order touches order rows, line items and stock, which is a multi-row invariant. Nothing here rules out the default; the invariant actively argues for it.

The revenue report is the exception. It scans years of data and would compete with live traffic. Serve it from a read replica at first, and move it to a columnar store fed by change data capture when it grows. That is a second store justified by a measured need, not by fashion.

Contrast an IoT workload: a million devices each sending a 200-byte reading every 10 seconds is 100,000 writes per second, about 20 MB per second or roughly 1.7 TB per day before replication. Reads are almost all by device and time range, old data expires, and there are no multi-row invariants. Every question now points away from a single relational primary and towards a time-series or wide-column store with TTL-based expiry. The method is the same; the numbers decide.

Consistency and transactions, concretely

Ask what goes wrong if two requests race. If the answer is a duplicate notification, eventual consistency is fine. If the answer is selling the same seat twice, you need the store to arbitrate: a transaction, a unique constraint, or a conditional write. Key-value stores often provide conditional writes and limited transactions, which are enough when the invariant lives inside one item or a small fixed group of items. Invariants across arbitrary rows, or reports that must match the ledger, favour a relational or distributed SQL store.

Also decide what reads must see. Reading from replicas or eventually consistent indexes means a user may not see their own write. That is fine for a feed and confusing for a checkout page. Write the rule down per query in the inventory.

Benchmark the shortlist properly

Benchmarks published by vendors measure the vendor's best case. Yours must measure your case: your schema, your query mix in the right proportions, realistic data volumes, and at least twice your forecast peak. Run long enough to see compaction, vacuum, cache churn and checkpoints, because p99 latency after an hour often differs from the first five minutes. Then break things: kill the primary or a node during the run and record recovery time and the errors clients saw.

# Benchmark the candidate with YOUR shapes, at 2x peak, and watch p99 over time.
pgbench -i -s 200 bench                       # or load a scrubbed copy of real data
pgbench -c 64 -j 8 -T 1800 -P 10 -R 2300 \
        -f get_order.sql@8 -f orders_by_customer.sql@4 -f create_order.sql@1 \
        -f update_status.sql@2 bench
# During the run: kill the primary, measure time to recovery and errors seen by clients.

For non-relational candidates, YCSB or a small purpose-built load generator plays the same role. The output that matters is a table of p50 and p99 per query type at 1x and 2x peak, plus failover behaviour, for each candidate.

Operations, cost and lock-in

  • Backup and restore. Point-in-time recovery and a timed restore drill. A backup you have never restored is a hypothesis.
  • Schema and data changes. How do you add a column or backfill a billion rows online? Some stores make this easy; key-value stores turn it into a data migration you write yourself.
  • Upgrades and failover. Major-version upgrades, replica promotion and connection handling during failover all need runbooks.
  • Cost shape. Managed key-value services often bill per request, so cost tracks traffic; provisioned clusters cost the same idle or busy. Model both at today's load and at 10x.
  • Lock-in. Proprietary query models and APIs make leaving expensive. Fine if you choose it knowingly; write down the exit cost.

When a second store is justified

Polyglot persistence is powerful and expensive. Each extra store adds a consistency boundary, an on-call surface and a sync pipeline. Add one when a measured need cannot be met by the primary: search relevance, large analytical scans, very high-rate telemetry, similarity search. Feed it from the system of record with change data capture or an outbox table, never with dual writes from the application. Dual writes drift as soon as one of the two writes fails. If the pipeline needs a broker in between, How to pick a message broker applies the same workload-first method.

Failure modes

  • Choosing for scale you do not have. A distributed store chosen for a workload that fits on one node costs you joins, transactions and operational simplicity for nothing.
  • Modelling entities in a query-first store. In wide-column and key-value stores you design tables around queries. Copying a relational schema leads to scans, hot partitions and a redesign.
  • Hot keys. Totals look fine while one tenant, device or date partition takes most of the traffic. Put key skew in the inventory.
  • The database as a queue. Polling tables for work at high rates causes lock contention and bloat; use a queue or a proper outbox pattern.
  • Ignoring the analytics path. Heavy reports on the primary degrade production. Plan the replica or analytical copy early.
  • No restore drill. The first restore attempt happens during an incident, and it fails.

Trade-offs

There is no best database, only a best fit for a workload and a team. Relational stores trade some write scalability for flexibility and safety. Key-value and wide-column stores trade query flexibility for predictable scale. Managed services trade control and sometimes cost for less operational work. Specialised stores trade an extra moving part for one capability done well. The method keeps these trades explicit: every choice should point to a line in the query inventory or a number in the sizing that demanded it.

What to do next

  1. Write the query inventory for your system: every read and write, peak rate, latency objective and lookup key.
  2. List every multi-row invariant and decide how the store must enforce each one.
  3. Run the sizing arithmetic with a growth margin, and note key skew.
  4. Start from the relational default and write down any number that rules it out.
  5. Benchmark at most three candidates with your own query mix at 2x peak, including a failover during the run.
  6. Decide the analytics and search paths now, fed by CDC or an outbox, not dual writes.
  7. Schedule a restore drill before launch and record the time it took.
Key takeaway: Choosing a database is a workload problem, not a product problem. Write down every query with its rate, latency objective and key, list the invariants that span rows, size the data and traffic with a growth margin, and let those numbers eliminate options. A relational database is the right default until a measured requirement rules it out. Benchmark the shortlist on your own query mix, including failure, and add specialised stores only for needs the primary cannot meet, fed from the system of record rather than by dual writes.