Database Selection and Trade-offs Questions
Choosing the right database and data platform for a workload: relational versus NoSQL versus specialized stores, managed versus self-hosted, and matching technology to consistency, scale, cost, and query or access-pattern needs. Covers OLTP versus OLAP and transactional-versus-analytical workload splits, polyglot persistence across multiple data stores, structuring an ambiguous selection prompt, naming trade-offs, and defending a recommendation to stakeholders.
You must pick a primary datastore for a payments system requiring strict transactional consistency, fraud detection analytics, and monthly reconciliation reports. Propose an architecture that may include more than one data platform, explain data flow between components, and justify choices based on consistency, latency, and analytics needs.
Sample Answer
Direct answer
For a payments system, split the workload across two purpose-built platforms and connect them with a change stream, rather than forcing one database to be good at everything. Use a strongly consistent relational database (or a distributed SQL system like Google Cloud Spanner) as the single system of record for the money-moving path, stream every committed change out via change data capture (CDC, a mechanism that reads a database's internal transaction log and turns each row change into an event), and land that stream in a columnar analytics warehouse for fraud scoring and monthly reconciliation. The OLTP database (online transaction processing: a system optimized for many small, concurrent reads/writes with strict correctness) never talks directly to the analytics side in the request path, so a slow fraud query or a big reconciliation scan can never add latency to a customer's payment.
Architecture and data flow
Component 1: system of record (OLTP). A relational database (PostgreSQL/MySQL-compatible, or a managed option like Amazon Aurora or Cloud Spanner) holds accounts, ledger entries, and transaction state. Every payment is a single ACID transaction (atomicity, consistency, isolation, durability: the guarantee that a transaction either fully applies or not at all, leaves data valid, is isolated from concurrent transactions, and survives a crash once committed), run at SERIALIZABLE or the engine's strictest practical isolation level, so two concurrent debits against the same balance can never both succeed and double-spend the account. This is the only place writes to money happen.
Component 2: change stream. A CDC connector (Debezium is the common open-source choice) tails the OLTP database's write-ahead log and publishes one event per committed row change onto a durable log such as Kafka. This decouples "the transaction committed" from "everyone who needs to know about it finds out," so downstream consumers can be slow, restarted, or temporarily down without blocking payments.
Component 3: real-time fraud scoring. A stream processor (Kafka Streams or Flink) consumes the change stream, joins it against recent history and device/session signals, and scores each transaction. High-risk transactions write a hold flag back to the OLTP system through the normal transactional API, idempotently keyed by transaction ID so a retried or duplicated event cannot apply the same hold twice.
Component 4: analytics and reconciliation. The same change stream is sinked into a columnar warehouse (Snowflake, BigQuery, or Redshift). Monthly reconciliation jobs run here: compare internal ledger totals against external bank/processor settlement files, using checksums and row counts per batch to prove nothing was lost or double-counted in transit. This workload can run full-table scans and heavy joins without any risk to the OLTP system, because it reads from a copy, not the source.
flowchart LR
Client[Payment client] --> API[Payment API service]
API -->|write ledger row, serializable txn| OLTP[(Relational OLTP\nsystem of record)]
OLTP -->|CDC stream| Bus[[Kafka change stream]]
Bus --> Fraud[Stream fraud scorer]
Fraud -->|hold flag, idempotent write-back| OLTP
Bus --> WH[(Columnar warehouse\nreconciliation and analytics)]
Bus --> Audit[(Append-only audit log)]
Worked example: why single-platform doesn't fit
Say the service processes 500 transactions per second (TPS) at peak, and a monthly reconciliation job needs to scan 90 days of transaction history (roughly 500 * 86,400 * 90 ≈ 3.9 billion rows) and join it against a settlement file. Running that scan directly against the OLTP primary would compete for the same buffer pool (the in-memory area where the database engine caches recently used data pages so it doesn't have to re-read them from disk) and lock manager that the live payment path depends on; even a well-indexed OLAP-style (online analytical processing: the counterpart to OLTP, large scans and aggregations for reporting rather than many small transactional writes) query with a large sort or hash join can hold shared resources long enough to add tail latency to concurrent writes. Separating the two means the OLTP database is sized and tuned purely for point lookups and small transactional writes (small working set, aggressive caching, low lock contention), while the warehouse is sized purely for sequential scans and columnar compression over the same 3.9 billion rows, with no shared contention between them.
Trade-offs and pitfalls
- Consistency boundary matters. The fraud-hold write-back has to be idempotent and the CDC pipeline has to guarantee at-least-once delivery with a de-duplication key downstream, or a redelivered event could double-apply a hold or a reversal. "Exactly-once" in practice means at-least-once delivery plus an idempotent consumer, not a magic delivery guarantee.
- Operational cost. This design adds CDC, a message bus, a stream processor, and a warehouse: four extra systems to monitor, upgrade, and page on. That overhead is justified here because payments cannot tolerate the alternative (a single database straining to be both a fast ledger and a fraud/analytics engine), but it would be over-engineering for a system without the same consistency and scale pressure.
- Reconciliation needs a proof, not just a copy. A common mistake is treating the warehouse copy as automatically correct. Reconciliation should independently verify completeness (row counts, checksums, watermarking on the CDC stream) rather than assuming replication never drops or duplicates an event.
- **A single globally-consistent NewSQL system (NewSQL: a class of databases that aim to keep relational ACID guarantees while scaling out like a distributed system, e.g. Spanner or CockroachDB) can replace the OLTP component if the team wants one fewer moving part for the transactional side, but it does not remove the need for a separate analytics store: OLAP-shaped queries still shouldn't run against the transactional cluster.
You are evaluating NewSQL systems (CockroachDB, Google Spanner) vs sharded PostgreSQL for an application needing serializable isolation and global scale. From an SRE standpoint, compare operational complexity, latency, cost, schema migrations, backup/restore, and failure recovery modes.
Sample Answer
Direct answer
Recommend NewSQL (Cloud Spanner or CockroachDB) over sharded PostgreSQL for a workload that genuinely needs serializable isolation at global scale (the strongest transaction-isolation level, guaranteeing that concurrent transactions behave as if they had run one at a time in some order, with no partial or conflicting reads and writes visible in between), because sharded PostgreSQL cannot give that guarantee across shards without the team building its own distributed-transaction coordination on top of it. Placed on an online transaction processing (OLTP) latency spectrum alongside the workload's other realistic options: vertically-scaled PostgreSQL has the lowest latency but a hard write-scaling ceiling; Aurora improves storage, read scaling, and availability over vanilla PostgreSQL but is still a single-writer system, not a horizontal write-scaling answer; CockroachDB gives horizontal write scaling and true cross-shard serializable transactions at the cost of consensus round-trip latency; Cassandra gives the best raw write throughput and lowest latency of the four but does not offer multi-key serializable transactions at all, which is why it is excluded from the shortlist even though it is the fastest option in the group.
Structured elaboration
Why sharded PostgreSQL fails the stated requirement. Sharding PostgreSQL, whether through an extension like Citus or through hand-rolled application-level routing, distributes data across independent PostgreSQL instances. A transaction that only touches rows on one shard is a normal, fully serializable PostgreSQL transaction. A transaction that touches rows on two different shards is not, unless the team builds a distributed-transaction protocol (typically a two-phase commit, 2PC, coordinating a commit decision across multiple independent databases) on top of it themselves, and even a correctly built 2PC layer does not automatically give the same serializability guarantee a native distributed database provides across arbitrary key ranges. NewSQL systems solve this as a first-class, built-in primitive: both Spanner and CockroachDB implement genuinely distributed serializable transactions spanning any set of keys, anywhere in the cluster, without the team building that machinery.
The four-way OLTP latency comparison, named explicitly.
| System | Write scaling | Typical latency profile | Cross-row transaction guarantee |
|---|---|---|---|
| Vertically-scaled PostgreSQL | None past the largest single instance available | Lowest, no distributed-consensus overhead | Full serializability, but only within one instance |
| Aurora | Still single-writer; Aurora scales storage durability, read replicas, and failover, not write throughput | Low, close to vanilla PostgreSQL, slightly better tail latency from its distributed storage layer | Full serializability, but still only within one writer |
| CockroachDB | Horizontal, writes distribute across the cluster | Higher than the two above, a consensus round trip (Raft) is paid per transaction, more again across regions | Full distributed serializability across arbitrary shards |
| Cassandra | Horizontal, the best raw write throughput of the four | Lowest of the distributed options, no cross-partition transaction coordination to pay for | None: no multi-key serializable transactions at all |
The ordering matters for the recommendation: Cassandra is the fastest and highest-throughput of the four, and it is excluded specifically and only because it does not meet the isolation requirement the question states, not because it performs poorly. Aurora is a genuine improvement over vanilla PostgreSQL for availability and read scaling, and it is excluded from being the write-scaling answer specifically because its writer remains a single instance; conflating "Aurora scales" with "Aurora scales writes horizontally" is a common and costly misreading of what Aurora actually does. One caveat to that framing: Amazon Aurora Limitless Database, now generally available for Aurora PostgreSQL, does provide horizontal write scaling by automatically sharding data and queries across multiple underlying Aurora Serverless instances while preserving the transactional consistency of a single database, so "Aurora can never scale writes horizontally" is no longer true of the Aurora product line as a whole. It is still a narrower answer than Spanner or CockroachDB for this question specifically: as of this writing it operates within a single AWS Region rather than synchronously across regions, so it does not meet a genuinely global-scale requirement on its own, which is why the comparison and recommendation above still hold for this question's global-scale shortlist, but a team facing a single-region horizontal-write-scaling need should check Limitless Database before assuming Aurora is disqualified.
Operational complexity from an SRE standpoint. A fully managed NewSQL option (Spanner, or CockroachDB Cloud) has the lowest ongoing operational burden: no consensus infrastructure, rebalancing, or certificate management to run by hand. Self-hosted CockroachDB has the highest burden of the distributed options, since the team runs that infrastructure itself. Sharded PostgreSQL sits in an awkward middle: each individual shard is operationally familiar (standard PostgreSQL), but the shard-routing and resharding logic is bespoke, and every schema migration has to be coordinated by hand across every shard, since there is no native cluster-wide distributed data-definition-language (DDL) operation the way NewSQL systems provide.
Backup and restore. NewSQL systems support consistent, cluster-wide, point-in-time backups (Spanner's managed backups; CockroachDB's BACKUP/RESTORE with as-of-system-time consistency across the whole distributed cluster). Sharded PostgreSQL needs per-shard backups coordinated to a consistent point in time by the team's own tooling, and restoring a cross-shard-consistent snapshot is meaningfully harder to guarantee correctly.
Failure recovery. NewSQL systems self-heal a single-node or single-availability-zone loss automatically through consensus re-election, with no human intervention. A sharded PostgreSQL failover (promoting a replica to take over as the primary when the existing primary becomes unavailable) for one shard is a well-understood, standard PostgreSQL failover (commonly automated with a tool like Patroni), but it does not automatically preserve cross-shard consistency, and losing one shard's availability zone leaves that shard's rows unavailable while the rest of the system continues, a partial outage some teams consider acceptable and others consider a hidden complexity they did not sign up for.
Cost. This is the one axis where sharded PostgreSQL can genuinely win at moderate scale: NewSQL pays a real "consensus tax," every write commits only after a quorum of replicas (commonly three) acknowledges it, meaning at least one extra network hop versus a single-writer system, and the storage and compute footprint is multiplied by the replication factor. Sharded PostgreSQL, without that tax, can be the cheaper option per unit of throughput at a scale where its operational and correctness gaps have not yet become a real problem.
Worked example
Size a NewSQL cluster against a stated sustained transaction rate, to make the "consensus tax" concrete.
# NewSQL node sizing for a global-scale OLTP ledger.
import math
target_tps = 3000
node_tps_capacity = 800 # illustrative per-node sustained write throughput in a synchronously-replicated cluster
nodes_needed = math.ceil(target_tps / node_tps_capacity)
replication_factor = 3 # quorum-based replication, a common default
raw_write_amplification = replication_factor # each logical write becomes this many physical writes across replicas
print(f"target sustained TPS = {target_tps}")
print(f"nodes needed at ~{node_tps_capacity} TPS/node = {nodes_needed}")
print(f"physical write amplification from replication factor {replication_factor} = {raw_write_amplification}x")
# target sustained TPS = 3000
# nodes needed at ~800 TPS/node = 4
# physical write amplification from replication factor 3 = 3x
Four nodes to sustain 3,000 transactions per second sounds modest, but every one of those logical writes is really three physical writes under the hood, once for each replica in the quorum. That is the concrete shape of the "consensus tax": it shows up as both the extra network round trip per write (the latency cost) and roughly a tripling of raw storage and compute footprint relative to a single-writer system handling the same logical write rate (the cost line item), and it is the real number to weigh against sharded PostgreSQL's lower per-unit cost when scale has not yet forced the correctness question.
Trade-offs & pitfalls
- Conflating "Aurora scales" with "Aurora scales writes horizontally." Standard Aurora is a strong choice for read scaling, durability, and failover, and its default configuration does not remove the single-writer ceiling; recommending vanilla Aurora as the answer to a horizontal-write-scaling requirement is a real, recurring mistake. Amazon Aurora Limitless Database is the documented exception (single-Region horizontal write scaling for Aurora PostgreSQL), so name it explicitly rather than treating "Aurora" as monolithically single-writer.
- Excluding Cassandra from a serializable-transaction requirement while still praising its raw performance out of context. Its exclusion here is entirely about the missing correctness guarantee, not a performance judgment; naming that distinction explicitly avoids implying Cassandra is simply "worse."
- Building a cross-shard transaction layer on sharded PostgreSQL by hand and assuming it now provides the same guarantee a native NewSQL system does. A homegrown 2PC layer is real engineering effort with its own failure modes, and it is easy to overestimate how close it gets to genuine distributed serializability.
- Underestimating the schema-migration coordination cost of sharded PostgreSQL until the first migration that has to run consistently across dozens of shards without a native distributed DDL mechanism.
- Choosing NewSQL by default "for correctness" without checking whether the workload's actual transaction rate and region count justify the consensus tax. At small to moderate single-region scale, the tax may not be worth paying yet.
You are evaluating a storage layer for an online retail product catalog that must support rich queries (joins, filters), transactional updates (price changes, stock), and occasional schema changes (new attributes). Compare relational databases vs document NoSQL stores for this use case. Describe the trade-offs around schema flexibility, transactional guarantees, query expressiveness, scaling patterns, and operational complexity. Recommend a choice and justify under what conditions you'd pick the other option.
Sample Answer
Direct answer
Default to a relational database, for example PostgreSQL, for this catalog. The three stated needs, rich queries with joins and filters, transactional updates to price and stock, and occasional schema changes, all point toward relational's strengths, not a document store's. Joins and filters are relational's core capability, transactional updates need ACID guarantees (atomicity, consistency, isolation, durability, the standard relational correctness properties) rather than eventual consistency, and "occasional" schema changes are exactly what a modern relational engine handles well, either through a normal migration or a JSONB column for the genuinely variable attributes.
Structured elaboration
Schema flexibility. "Occasional new attributes" is a weak reason on its own to pick a document store over relational. Adding a nullable column to a large table is typically a fast, metadata-only operation on modern PostgreSQL or MySQL; only rewriting an existing large table under a non-null default with no default value is genuinely expensive. For attributes that truly vary per product category, a shirt has size and color, a laptop has RAM and storage, a JSONB column inside the relational schema captures that document-like flexibility for just that sub-shape, while price, stock, and category stay relational and joinable.
Transactional guarantees. Price changes and stock decrements are exactly the "update several related facts, all or nothing" shape relational transactions exist for. A document store's atomicity usually covers a single document's own fields well, but a promotion touching many products at once, apply twenty percent off to an entire category, right now, or not at all, needs cross-document transactional guarantees. Some document databases do support multi-document transactions, but at a real performance cost relative to their native single-document write path, versus relational where that is the native case.
Query expressiveness. The question names joins and filters directly, and joins are relational's headline strength. A document approach models relationships either by embedding related data into the document, fast reads but now a price update must be propagated to every embedded copy, or by referencing an identifier, which the store will not join for you, so the application performs repeated lookups or builds a separate aggregation pipeline. Neither is as direct as a relational join for a previously unplanned query shape.
Worked example
Model Product(id, sku, price, stock_qty, category_id) normalized against Category, with a flexible attributes JSONB column carrying category-specific fields, queryable via a GIN index (a Generalized Inverted Index, PostgreSQL's index type for searching inside a semi-structured value like JSONB rather than a plain scalar column) for the rarer cases that do need to filter on a dynamic attribute. Contrast that with a pure document approach, where category-specific fields embed naturally per document, but a "price update across an entire category" operation becomes either many individual document updates with no cross-document atomicity by default, or a multi-document transaction the engine treats as a heavier operation than its normal write path.
Trade-offs and pitfalls
A pitfall is choosing a document store purely because "the catalog has flexible attributes," without checking whether the flexible part is a small subset of the schema, the attributes, rather than the whole schema, price, stock, and category, which stay uniform. A JSONB column captures exactly that shape inside relational without giving up joins or transactions on everything else.
A second pitfall is assuming relational cannot scale reads. For a catalog, read scaling is normally solved with read replicas and a cache in front of product pages, not by switching the underlying data model.
What would flip the recommendation toward a document store: a marketplace where every seller or category defines its own attribute schema with almost no shared structure across categories, thousands of distinct and evolving shapes, not "occasional" changes to a broadly shared schema, and where queries are predominantly "fetch one product by identifier" or "list one seller's products" rather than cross-category filtering and joins. At that point a document store's per-document schema freedom outweighs relational's join strength, because there is comparatively little left to join.
You must design a global, low-latency user profile store for 200M users supporting reads under 20ms from any region and occasional writes (profile updates). Candidate technologies: DynamoDB Global Tables, Cassandra multi-dc, PostgreSQL with read replicas. Choose a solution, justify it, and outline replication, conflict resolution, cost implications, and read routing.
Sample Answer
Direct answer
Recommend DynamoDB Global Tables for the profile-store use case as stated: writes are occasional and profile fields rarely conflict across regions at the same instant, so eventually-consistent, last-writer-wins (LWW) cross-region replication is an acceptable trade for fully managed, active-active writes with local reads under 20ms in every region.
Structured elaboration
Why every candidate needs a LOCAL copy of the data in each read region. A synchronous cross-region round trip alone typically busts a 20ms budget by an order of magnitude. The physics floor for a round trip is the distance divided by the speed of light in fiber (roughly two-thirds the speed of light in a vacuum, due to the refractive index of glass):
import math
def great_circle_km(lat1, lon1, lat2, lon2, R=6371):
lat1, lon1, lat2, lon2 = map(math.radians, [lat1, lon1, lat2, lon2])
dlat, dlon = lat2 - lat1, lon2 - lon1
a = math.sin(dlat/2)**2 + math.cos(lat1)*math.cos(lat2)*math.sin(dlon/2)**2
return R * 2 * math.asin(math.sqrt(a))
d_km = great_circle_km(39.0, -77.5, 1.35, 103.8) # N. Virginia <-> Singapore
c_fiber_km_s = 300_000 * (2/3)
rtt_ms = (d_km / c_fiber_km_s) * 2 * 1000
print(f"great-circle distance: {d_km:,.0f} km")
print(f"physics-floor RTT (straight-line fiber, no routing overhead): {rtt_ms:.1f} ms")
great-circle distance: 15,526 km
physics-floor RTT (straight-line fiber, no routing overhead): 155.3 ms
Real internet paths are always slower than this straight-line floor (routing hops, congestion), so a design that answers every read with a synchronous cross-region call cannot meet a 20ms budget at all; every candidate must serve reads from a copy of the data that already lives in the requesting region.
"Occasional writes" does not demand low write latency from every region, only that writes eventually land everywhere without being lost.
Candidates.
- DynamoDB Global Tables: replicates via DynamoDB Streams to every region; every region can read AND write locally; conflicts are resolved last-writer-wins (LWW) by internal timestamp; propagation between regions is typically sub-second, though AWS documents no hard SLA on that figure; fully managed, billed per region's capacity plus a replicated-write-unit charge for each additional region.
- Cassandra multi-DC:
NetworkTopologyStrategyplaces replicas per datacenter, andLOCAL_QUORUMreads/writes stay in-region for low latency while replication to other datacenters happens in the background, a similar shape to Global Tables but with tunable, per-request consistency rather than a single fixed LWW rule. This is the natural pick if the organization is not committed to a single cloud provider: Cassandra or Scylla multi-DC is the closest cross-cloud equivalent to what Global Tables gives on AWS alone, at the cost of self-managing the cluster (or paying a managed-Cassandra vendor). - PostgreSQL with read replicas: one write-primary region, asynchronous streaming replication to per-region read replicas. Reads are local and fast; but EVERY write funnels through the one primary region regardless of where the user is, so a write from a non-primary region pays the full cross-region round trip (tolerable for genuinely occasional profile edits, a poor fit the moment writes need to feel local too). It is also a single point of write failure: if the primary region goes down, every other region loses write capability until a primary is promoted, and any replica that had not caught up loses its most recent writes on promotion.
Read routing. Route each request to the nearest region's local replica, for example with latency-based DNS routing (directing a client to whichever region answers fastest) or an edge/CDN layer that forwards to the nearest regional API, so the read never leaves the region.
flowchart LR
U1[User: EU] --> R1[EU region: local Global Table replica]
U2[User: US] --> R2[US region: local Global Table replica]
U3[User: APAC] --> R3[APAC region: local Global Table replica]
R1 <-. async replication .-> R2
R2 <-. async replication .-> R3
R1 <-. async replication .-> R3
Worked example
Two regions each write the same user's display_name field within the same short replication window, close enough together that both writes are in flight before either has propagated. DynamoDB Global Tables resolves the conflict by keeping whichever write has the later internal timestamp and silently discarding the other, no error, no merge. The practical implication: don't build a "who last touched this field" audit feature purely off the base table, since the losing write leaves no trace there; log profile changes to a separate append-only table or stream if that provenance genuinely matters.
Trade-offs & pitfalls
What would flip the recommendation. For a REAL-TIME game-data variant of this same shape (position, health, inventory, written far more frequently, usually from a session pinned to one region at a time), the calculus changes: Cassandra/Scylla multi-DC with LOCAL_QUORUM (tunable, and not tied to one vendor's specific LWW rule) is the better fit, especially on a multi-cloud platform. If the organization is not on AWS at all, Cassandra/Scylla multi-DC answers that constraint directly, since it runs the same way on any cloud or on-premises. Two pitfalls to name explicitly: Global Tables' replicated-write-unit cost scales with the number of regions (each additional region adds roughly another full write's worth of cost), and "no SLA on propagation" means there is no contractual recovery-point guarantee, only an empirically sub-second one; anything that truly cannot tolerate an unbounded (if rare) propagation delay needs a design that does not rely on Global Tables' default consistency model at all.
You must evaluate three candidate databases for a write-heavy leaderboard system: Redis (in-memory), Cassandra (wide-column), and PostgreSQL (disk-backed). Define benchmark scenarios (writes/sec, reads/sec, read-after-write latency, data size), key failure modes to test, and what metrics and SLOs you would use to pick the right platform.
Sample Answer
Direct answer
Expect Redis to win a leaderboard-shaped write-heavy benchmark, because its sorted-set data structure is purpose-built for exactly this access pattern, with Cassandra as the fallback once the working set genuinely outgrows what fits comfortably in memory, and PostgreSQL as the control that will most likely lose on raw write throughput but establishes the transactional-safety baseline the other two are measured against. That expectation is a hypothesis, not the answer: the actual answer to "evaluate three candidates" is the benchmark design below, built so the result is decided by measurement, not by which system sounds best on paper.
Structured elaboration
Benchmark scenarios to define, each with an explicit target, not a vague description.
| Dimension | What to define | Illustrative target for this exercise |
|---|---|---|
| Writes per second | Sustained score-update rate during a realistic peak (a live event, a tournament) | See worked example: derived from a stated player and update-rate assumption |
| Reads per second | Leaderboard views (top-N plus "my rank") during the same peak | Derived from a stated viewer-to-player ratio |
| Read-after-write latency | Time from a score update to that update appearing in a subsequent top-N read | Under 1 second at the 99th percentile (p99), a reasonable real-time-leaderboard expectation |
| Data size | Number of distinct players in the working set, and average bytes per entry | Drives whether the working set fits the tested memory budget, the central question for Redis specifically |
Failure modes to test, one per candidate plus one cross-cutting. Node or primary failure mid-write-burst: for Redis, verify whether a failover (the automatic process of promoting a standby replica to primary and redirecting traffic when the current primary goes down) during an unacknowledged asynchronous replication window can lose the most recent writes, and whether that risk is acceptable or needs to be closed with a stronger write-acknowledgment setting (accepting the latency cost of waiting for a replica to confirm). For Cassandra, verify behavior under a network partition: does it continue accepting writes per the configured consistency level, and do replicas correctly reconcile once the partition heals. For PostgreSQL, verify failover behavior with a managed or self-managed replication topology and measure how long writes are actually unavailable during the failover window. Cross-cutting: run a genuine network partition or node-kill during the sustained peak load from the writes/reads targets above, not against an idle system, since failure behavior under load is the behavior that matters.
Metrics and service-level objectives (SLOs) to observe. Write latency, p50 and p99. Read latency, p50 and p99. Read-after-write staleness distribution (not just an average, the tail is what a user actually notices). Replication lag under load. Error rate during the fault-injection window. Time to return to SLO-compliant behavior after a failure resolves. And one candidate-specific early-warning metric each: for Redis, memory headroom against maxmemory and the active eviction policy (if the working set exceeds available memory, Redis evicts keys per policy, commonly least-recently-used, which can silently and incorrectly drop leaderboard entries if the eviction policy is not deliberately scoped away from leaderboard keys); for Cassandra, compaction backlog (a growing backlog under sustained high write rate is the leading indicator that read latency is about to degrade, well before it visibly does); for PostgreSQL, table and index bloat ratio from autovacuum lag (a write-heavy pattern that repeatedly updates the same rows, exactly what leaderboard score updates do, is the classic case that causes PostgreSQL's multi-version concurrency control, MVCC, to accumulate dead row versions faster than autovacuum reclaims them, degrading both read and write latency over time if unaddressed).
Why these three failure modes are not interchangeable, and why they matter more than the raw throughput numbers. A benchmark that only measures steady-state throughput will make all three candidates look reasonable; the failure modes above are each a specific, well-known weak point of a write-heavy pattern on that particular engine, and a benchmark that skips them will not catch the actual production incident each engine is prone to.
Worked example
Derive the writes/sec and reads/sec targets from a stated scenario, and check Redis's working-set memory budget against a stated player count, since memory headroom is the single most consequential capacity question for a Redis-based design.
# benchmark target derivation for a write-heavy leaderboard.
players = 5_000_000
updates_per_player_per_day = 40 # match/score events per active player per day
writes_per_day = players * updates_per_player_per_day
writes_per_sec_avg = writes_per_day / 86400
peak_multiplier = 6 # evening/event-driven peak concentration
writes_per_sec_peak = writes_per_sec_avg * peak_multiplier
print(f"avg writes/sec = {writes_per_sec_avg:,.0f}, peak writes/sec (x{peak_multiplier}) = {writes_per_sec_peak:,.0f}")
# Redis working-set sizing for a large leaderboard.
lb_players = 20_000_000
bytes_per_entry = 200 # member id + score + sorted-set node overhead, illustrative
working_set_gb = lb_players * bytes_per_entry / 1e9
replication_factor = 2 # primary + 1 replica for HA
headroom = 1.3
total_ram_gb = working_set_gb * replication_factor * headroom
print(f"working set = {working_set_gb:.1f} GB for {lb_players:,} players")
print(f"RAM budget with {replication_factor}x replication and {headroom}x headroom = {total_ram_gb:.1f} GB")
# avg writes/sec = 2,315, peak writes/sec (x6) = 13,889
# working set = 4.0 GB for 20,000,000 players
# RAM budget with 2x replication and 1.3x headroom = 10.4 GB
Two results worth acting on. First, peak write load (about 13,900 writes per second under these stated assumptions) is well within what all three candidates can sustain in isolation; the benchmark's value is in the failure-mode and tail-latency behavior at that load, not in whether any of the three can technically keep up. Second, a 20-million-player leaderboard's working set is only about 4 gigabytes, comfortably fitting in a single modestly sized Redis instance with room to spare for replication and headroom, well under commonly available managed-instance memory tiers. The real trigger for moving off Redis toward Cassandra is not this leaderboard's size, it is a leaderboard an order of magnitude or two larger, or one that needs to retain full historical score events (not just current standings) indefinitely, which no longer fits an in-memory design economically.
Trade-offs & pitfalls
- Benchmarking only steady-state throughput. All three candidates will look acceptable; the failure modes above are where the real differentiation, and the real production risk, actually lives.
- Assuming Redis's async replication failover is safe by default. It is not, without an explicit stronger write-acknowledgment configuration, which trades some write latency for the durability a leaderboard's ranking correctness actually needs.
- Missing a Cassandra compaction backlog until read latency has already visibly degraded. Compaction backlog is a leading indicator specifically because it degrades before the symptom (slow reads) becomes obvious; alert on the backlog metric itself, not just on read latency.
- Treating PostgreSQL's bloat risk as a generic "Postgres is slower" conclusion, rather than the specific, addressable cause it is: hot-row updates outrunning
autovacuum. Tuningautovacuumaggressiveness for the specific hot table, or restructuring the update pattern, is a real fix, not a reason to dismiss PostgreSQL outright for smaller-scale versions of this workload. - Setting the read-after-write latency SLO too strictly for the actual product requirement. A live global leaderboard tolerating roughly a second of staleness is normal and expected; over-specifying "instant" consistency here adds real engineering cost for a UX improvement users are unlikely to notice.
Unlock Full Question Bank
Get access to all Database Selection and Trade-offs interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.