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.
Your team must choose between single-region managed RDBMS with strong consistency and a globally replicated NoSQL store with eventual consistency. The service stores user shopping carts and must guarantee that an item added is visible when the user immediately fetches their cart. Evaluate the correctness and UX implications of both choices and propose a design that preserves user expectations across geographic regions.
Sample Answer
Direct answer
The requirement, an item added must be visible the instant the same user fetches their cart, is a read-your-writes guarantee, not a demand for global strong consistency at all times. Read-your-writes only needs to hold for a user relative to their own prior write, so the right design is not a binary pick between one strongly consistent relational database and one globally replicated eventually consistent store. It is: route each user's cart traffic to a single home region, get strong consistency cheaply within that region, and replicate asynchronously across regions for durability and for the rare case the user changes location.
Structured elaboration
Why a single-region strong RDBMS alone falls short. It trivially satisfies read-your-writes for users near that one region, but every other user pays cross-region latency on every cart operation, and that one region becomes a single point of failure for carts everywhere, which defeats the point of "globally replicated" in the first place.
Why a globally replicated NoSQL store's default mode alone falls short. A globally replicated key-value store's default multi-region mode typically replicates asynchronously with eventual consistency across regions: a write accepted in one region is not guaranteed to be visible to a read served from a different region within any bounded time. If a user's cart read happens to land on a different regional replica than the one that accepted their add-to-cart write, for example because a load balancer failed over between edge points of presence without session affinity, they can see a cart missing the item they just added. That is the exact correctness bug the question is pointing at.
The hybrid design.
- Home-region routing: assign each user or session a home region (nearest to their first-touch location) and pin both their cart reads and writes to it at the edge, via a routing header or cookie the load balancer inspects, rather than sending requests to whichever region happens to be least loaded.
- Same-region strong consistency for the common path: within the home region, either a normal relational transaction, or a strongly consistent read against the home region's replica of a key-value store, gives a real, cheap read-your-writes guarantee with no cross-region hop.
- Asynchronous cross-region replication for durability and travel: replicate cart data outward so it survives a regional outage and is still available, slightly stale, if the user changes location, at which point the system re-homes them to the new nearest region.
- For the narrow case where a user's device genuinely switches regions mid-session and home-region routing cannot be guaranteed, such as a server-to-server checkout integration without a sticky client, a synchronously replicated, strongly consistent write mode that some managed global key-value services now offer as an opt-in configuration can guarantee any-region reads see the latest write, at the cost of a real cross-region round trip on every write. That cost is acceptable selectively, not as the default for every add-to-cart click globally.
Worked example
Contrast the two failure paths concretely. A user in Sydney adds a jacket. If the client's next read is served by whatever region a stateless load balancer happens to pick, say Frankfurt because Sydney reported a momentary health blip, and Frankfurt has not yet received the asynchronous replica of that write, the fetch returns an empty cart. That is the literal bug. With home-region routing pinned to the Asia-Pacific region plus health-based regional failover, rather than per-request round robin across all regions, that scenario now only occurs during an actual regional outage of the user's home region, a much rarer and more acceptable failure mode than every single request carrying a chance of landing on a stale replica.
Trade-offs and pitfalls
A common pitfall is treating "eventual consistency" as a database property that can be swapped out later; the routing and session-affinity layer is doing the real work here, and changing the underlying database engine alone, relational to NoSQL or the reverse, does not fix a routing bug.
Another pitfall is defaulting to a synchronous, any-region-strong write mode everywhere because it sounds like the safest option. It taxes every single cart write with a real cross-region round trip, commonly tens of milliseconds depending on geography, for a guarantee the overwhelming majority of requests do not need, since most users never switch home regions mid-session.
What would flip the recommendation: the checkout and payment step, not the cart itself, needs a stronger guarantee independent of the cart's consistency model. Preventing a double charge when two payment attempts race from two devices calls for a genuinely serializable transaction (the strictest isolation level: the database guarantees the result is the same as if every transaction ran one at a time, with none interleaved) and an idempotency key at the payment boundary, regardless of how the cart's own reads and writes are routed.
Design the data storage architecture for a social feed service with 10M users and 1B posts. Requirements: personalized feeds with p95 read <100ms, write throughput 50k posts/min, support full-text search over posts, and eventual consistency acceptable for feed freshness. Map query patterns (fanout-on-write vs fanout-on-read) to storage options (wide-column, key-value cache, search engine) and justify replication, caching, and consistency trade-offs.
Sample Answer
Direct answer
Use a hybrid fanout: fanout-on-write (push a new post into every follower's precomputed feed at write time) for the vast majority of accounts, and fanout-on-read (pull the author's recent posts at read time and merge them into the feed) for high-follower accounts. Pure fanout-on-write collapses under the write amplification (one write turning into many downstream writes, here one post turning into one inbox write per follower) of a handful of extremely popular accounts, and pure fanout-on-read makes every ordinary read expensive by fanning out a query per followed account at read time.
Structured elaboration
Query pattern to storage mapping.
- Post storage (durable, high write volume, simple key lookup by post id): a wide-column store (Cassandra/Scylla) with replication factor (RF) of 3, a typical durability/availability choice giving each row three replicas so the loss of one node does not lose data, because posts are write-once, read-many and don't need relational joins.
- Precomputed per-user feed ("inbox"): a key-value cache (Redis, or a wide-column "feed" table keyed by user id holding a sorted list of post ids), because a feed read needs to be a single fast lookup by user id, not a fan-out query computed at read time.
- Full-text search over posts: a dedicated search engine (Elasticsearch/OpenSearch), a genuinely different query shape ("rank posts matching these words") that belongs in neither of the above.
Fanout math (worked). 50,000 posts/min is about 833 posts/sec:
posts_per_min = 50_000
posts_per_sec = posts_per_min / 60
avg_followers = 200 # illustrative average, for order-of-magnitude reasoning
naive_fanout_writes_per_sec = posts_per_sec * avg_followers
celeb_followers = 5_000_000
celeb_single_post_fanout = celeb_followers # one post, fanned out synchronously
print(f"posts/sec: {posts_per_sec:.1f}")
print(f"naive fanout-on-write at {avg_followers} avg followers: {naive_fanout_writes_per_sec:,.0f} inbox writes/sec")
print(f"a single post from a {celeb_followers:,}-follower account: {celeb_single_post_fanout:,} inbox writes for ONE post")
posts/sec: 833.3
naive fanout-on-write at 200 avg followers: 166,667 inbox writes/sec
a single post from a 5,000,000-follower account: 5,000,000 inbox writes for ONE post
The steady-state average (about 167,000 inbox writes/sec at this illustrative follower count) is a large but achievable rate for a wide-column store or Redis cluster built for high write throughput. The problem is the TAIL: one post from a 5-million-follower account would generate 5 million inbox writes in one event, dwarfing the steady-state load. For accounts past a monitored follower-count threshold, don't fan out on write; merge their few recent posts into a requesting user's feed at READ time instead, cheap because there are few of them and they are easy to cache. This hybrid is the standard, well-documented pattern large-scale feed systems use for exactly this problem.
Consistency and caching. Eventual consistency is explicitly acceptable per the requirements, so post-create can return immediately while fanout happens asynchronously in a background worker, and Redis in front of the wide-column feed store absorbs the read-hot top-of-feed page (well inside the p95 under 100ms budget; a cache miss falling through to the wide-column store is still fast, single-digit to tens of milliseconds).
flowchart TD
W[New post] --> C[(Cassandra: posts, RF=3)]
W --> F{Author follower count}
F -- normal account --> FW[Async fanout worker]
FW --> Redis[(Redis / wide-column feed store: per-user inbox)]
F -- celebrity account --> Skip[Skip fanout-on-write]
W --> Search[(OpenSearch: async index)]
Reader[Feed read] --> Redis
Reader --> Merge[Merge in celebrity posts at read time]
Trade-offs & pitfalls
The celebrity threshold needs an explicit, monitored cutover, a follower-count boundary that routes an account from fanout-on-write to fanout-on-read, or the hot-write problem reappears in production the first time an ordinary user suddenly goes viral. Full-text search staleness (posts indexed asynchronously) means a just-posted message might not be searchable for a few seconds, acceptable under the stated eventual-consistency tolerance, but worth stating explicitly rather than assuming.
Your team is considering migrating from a monolithic Postgres instance to a distributed NewSQL database to handle scale. What criteria would you use to decide whether the migration is worth it, and what would a migration checklist look like that minimizes data loss and downtime during the cutover?
Sample Answer
Direct answer
Migrate from a Postgres monolith to a distributed NewSQL database (a class of databases that gives you the familiar SQL relational model but spreads storage and transactions across many machines instead of one, something a single-primary relational database cannot do on its own) only with concrete evidence that vertical scaling, read replicas, and partitioning have a visible ceiling for the actual workload, not a growth projection, since this is one of the highest-risk, highest-effort changes available to a data platform team. When the evidence supports it, cut over online and reversibly: a CDC-based replication bridge running in parallel with the old system, verified before any traffic depends on it, and a fast rollback path kept alive through a burn-in window.
Structured elaboration
Is the migration worth it: decision criteria
- Concrete scaling evidence: is the system actually hitting a wall vertical scaling, read replicas, and partitioning cannot solve, for example a single-writer throughput ceiling, or a working set that no longer fits a single node's memory and I/O budget even after partitioning, not a projection of future growth.
- Multi-region write requirement: does the application genuinely need low-latency writes from multiple geographic regions with strong consistency, something a single-primary Postgres cannot do natively. This is the strongest legitimate driver for NewSQL specifically, as opposed to "more scale" in general, which sharded Postgres or a good caching layer often solves more cheaply.
- Team readiness: does the team have, or can it build within a reasonable time, real operational expertise in the target system's specific failure modes, which differ from Postgres's.
- Application compatibility: how much of the app's SQL, transaction shape, and driver usage is compatible with the target system's dialect and consistency model without a rewrite. Systems marketed as Postgres-compatible typically support a subset, with real gaps around cross-shard transactions and some SQL features.
- Cost at the new scale: is the new system's compute and licensing cost genuinely lower than the alternative of continuing to invest engineering time in scaling Postgres further.
A migration checklist that minimizes data loss and downtime
- Schema compatibility audit: run the actual DDL against the target system and catalog every incompatibility, data types, index types, extensions, stored procedures, before writing any migration code.
- CDC bridge: stand up a change-data-capture pipeline that streams every committed change from Postgres into the new system continuously, so the target stays a live, current replica rather than a one-time snapshot.
- Backfill and reconcile: bulk-load historical data while the CDC bridge runs, then reconcile, comparing row counts and checksums per table or shard between source and target, since the backfill and the live stream can race and only reconciliation proves they actually converged.
- Shadow reads: once reconciled, send a copy of read traffic to the new system without depending on its answers yet, and diff the results against Postgres's answers for the same queries, to catch correctness bugs before anything depends on them.
- Cutover with a fast rollback path: flip writes to the new system behind a routing layer that can be reverted in minutes, and keep the CDC bridge, or a reverse one, running through a defined burn-in window so a fallback to Postgres loses no data if something goes wrong.
- Decommission window: only after a burn-in period with no rollback triggered do you stop dual-writing and retire the old primary, keeping a final consistent snapshot of it for a defined retention period.
Worked example
Assume a 4 TB Postgres database with roughly 2 billion rows. Reconciliation checksums batches of 2,000 rows at a time; assume, as a stated planning assumption rather than a measured benchmark, a sustained checksum throughput of 5,000 rows per second for a single worker: 2,000,000,000 divided by 5,000 is 400,000 seconds, about 111.1 hours, roughly 4.6 days, if run serially. That is why reconciliation has to be parallelized across many workers in practice: splitting the same job across 50 parallel workers brings it to about 111.1 divided by 50, roughly 2.2 hours, a legitimate input for planning how long the reconciliation phase of a cutover window actually needs.
| Step | What it verifies | Rollback implication |
|---|---|---|
| Schema audit | The target can represent every table without silent data loss | Found early, costs nothing to reverse |
| Backfill + reconcile | Target's historical data matches source, byte for byte, not just row count | A mismatch here blocks cutover entirely, no traffic has moved yet |
| Shadow reads | Target's query results match source's for real production queries | A mismatch pauses cutover, no traffic depends on the target yet |
| Cutover | Writes succeed against the target under real load | Routing layer reverts to Postgres in minutes; CDC bridge kept running to avoid data loss during the revert |
| Decommission | New system has run cleanly through a full burn-in period | No longer reversible without restoring from the retained final snapshot |
Trade-offs & pitfalls
The single biggest risk is treating the backfill as the whole migration and skipping shadow reads: a schema or type-mapping bug, for example a NUMERIC column silently losing precision when mapped to the target's numeric type, often shows up only in query results, not in a row-count reconciliation. The single biggest operational mistake is decommissioning the old primary too early, before a full business cycle, a month-end close, for a financial workload, has run cleanly on the new system.
Describe workloads and trade-offs for key-value stores (Redis, DynamoDB) used as primary storage versus as a cache. Discuss persistence/durability options, eviction strategies, TTL usage, memory vs disk trade-offs, and scenarios where an in-memory KV store is acceptable as primary storage versus when persistence is required.
Sample Answer
Direct answer
A key-value (KV) store is safe to use as a system's PRIMARY, durable data source only when it durably persists every write to disk, with replication, before acknowledging it: DynamoDB, which synchronously replicates each write across multiple availability zones (AZs, physically separate data centers within a cloud region) before returning success, or Redis configured with append-only-file (AOF) persistence. Used purely as a CACHE in front of a durable primary store, an in-memory engine like Redis can run with persistence relaxed or off entirely, because a cache is disposable: it can always be repopulated from the source of truth after a restart.
Structured elaboration
Persistence and durability options
- Redis offers RDB (Redis Database) snapshotting: periodic, compact point-in-time dumps of the whole dataset, fast to restore from but able to lose everything written since the last snapshot on a crash. AOF (append-only file) logs every write command and replays it on restart; its
appendfsync everysecsetting (fsync to disk once a second) risks up to one second of data loss on a crash and is the default balance of speed and safety, whileappendfsync alwaysfsyncs on every write for near-zero loss at a real throughput cost. RDB and AOF can run together for both fast restarts and tight durability. - DynamoDB has no "durability off" switch: every write is synchronously written to storage and replicated across multiple AZs in the region before the write is acknowledged to the caller. That unconditional durability, not its key-value data model, is what makes it safe to use as a system of record.
Eviction strategies
Redis's maxmemory-policy setting decides what happens when memory fills: noeviction (reject new writes with an error instead of dropping data), allkeys-lru / allkeys-lfu (evict the least-recently-used or least-frequently-used key across the whole keyspace, ignoring TTLs), volatile-lru / volatile-lfu / volatile-random (only evict keys that have a TTL set, leaving permanent keys alone), and volatile-ttl (evict the key with the soonest expiry first). A CACHE almost always wants one of the allkeys-* or volatile-* policies, because losing a cached value is a performance hit, not a correctness bug. A PRIMARY store should run noeviction: silently dropping a "real" record because memory filled up is data loss, and an explicit write error the application can retry or alert on is the only acceptable failure mode for a system of record.
TTL (time-to-live) usage
DynamoDB TTL marks an item with an expiration timestamp; a background process deletes expired items, but AWS's own documentation says only that this typically happens "within a few days" of expiry, with no hard deadline, and it does not consume write capacity for the initial deletion. That makes DynamoDB TTL excellent for pruning things like expired shopping carts or old audit rows without a cron job, but it is not a real-time guarantee: never build a rate limiter or a lock timeout on the assumption that a DynamoDB TTL item disappears the instant it expires. Redis's EXPIRE is much closer to real time: it expires a key lazily on the next access and also via an active background sweep, so an expired session or lock is very unlikely to be visible past its TTL by more than a small window.
Memory vs disk trade-offs
Redis holds its whole working set in RAM, which gives sub-millisecond reads but makes capacity expensive (RAM costs far more per gigabyte than disk) and caps a single node's dataset size to what fits in memory. DynamoDB is disk-backed, so table size is effectively unbounded and storage is cheap by comparison, at the cost of a higher per-request latency floor (typically low single-digit milliseconds) than an in-memory hash lookup.
Worked example
Consider a login-session store for 5 million concurrently logged-in users, each session a roughly 2 KB blob (user id, roles, last-seen time, a CSRF token), with a 30-minute sliding expiry.
If a session is purely ephemeral (losing it just forces a re-login, nothing downstream depends on it surviving), it is safe to make Redis the PRIMARY store for it directly: SETEX session:<id> 1800 <blob> lets Redis's own TTL mechanism do the cleanup, volatile-ttl (or allkeys-lru if raw capacity pressure matters more than expiry order) is the right eviction policy, and AOF can reasonably be skipped or set to everysec, because "restart and force everyone to log in again" is an acceptable failure mode, not data loss.
sessions = 5_000_000
bytes_per_session = 2 * 1024 # ~2 KB blob
total_gb = sessions * bytes_per_session / (1024 ** 3)
print(f"{sessions:,} sessions x {bytes_per_session} bytes = {total_gb:.1f} GB working set")
5,000,000 sessions x 2048 bytes = 9.5 GB working set
That is a small, single-node-friendly Redis dataset. Now change one requirement: the session record must be retained for a compliance audit of "who was logged in as this session at time T", so it can never silently disappear on an eviction. That is no longer a cache, it is a durable record: put it in DynamoDB with a TTL attribute for eventual cleanup, and accept the "within a few days" deletion window, because nothing downstream needs millisecond-precise removal for an audit trail.
Trade-offs & pitfalls
The most common mistake is treating DynamoDB TTL as a deletion SLA: a team builds a distributed lock or a rate-limit window on "the item will be gone by the time I expect" and gets surprised by items lingering for hours or days; anything needing real-time expiry belongs in Redis's EXPIRE, not DynamoDB TTL. The second is leaving a Redis CACHE at its conservative noeviction default: once memory fills, writes start failing outright instead of quietly evicting, which can take down an application that assumed the cache would "just work." The third is a cost mistake in the other direction: sizing a Redis cluster for the entire dataset "just in case" a cold key gets requested, when only the hot working set needs to live in memory and everything else can fall through to the durable store on a miss.
Case study: an e-commerce business needs inventory availability checks with <50ms global latency for checkout and daily analytics for replenishment. Propose a storage architecture that serves fast reads for availability and eventual-consistency analytics. Include a migration plan from a single-region RDBMS to this architecture.
Sample Answer
Direct answer
Serve the checkout availability check from a fast, per-region cache backed by the existing relational database as the system of record for the actual commit, and feed daily replenishment analytics from an asynchronous export of that same relational data into a warehouse. Migrate incrementally, change data capture (CDC) feeding both new paths in parallel first, then cutting reads over region by region, rather than a single risky cutover, so the latency requirement and the analytics requirement are solved by two different, purpose-built read paths off one source of truth, not by one database trying to do both.
Structured elaboration
Why a single-region relational primary alone cannot hit 50 milliseconds globally. Even the fastest same-metro hop adds real latency, and a genuinely global user base means some checkouts happen physically far from wherever the one primary lives. Using an estimate of fiber propagation speed, roughly 68% of light speed in vacuum, about 203,900 km/s, a route of roughly 10,500 km, a rough distance between western Europe and Southeast Asia, gives a one-way propagation of 10,500 / 203,900 * 1000 ~= 51.5 ms, already over the entire 50 millisecond budget before any application logic runs, for users on the far side of that link.
The read path for checkout. Front each region with a fast, locally served read, an in-memory cache or a globally replicated key-value store's local-region read, holding current availability counts, refreshed continuously from the relational system of record via CDC. Checkout reads this local, fast, slightly-stale-by-design path for the "can this be added to the cart" check. The actual decrement at final purchase still goes through the relational primary's own transaction, or a regional compare-and-swap write with the primary as tie-breaker, so a brief staleness in the fast read path can never by itself cause an actual oversold unit.
The analytics path. Daily replenishment analytics needs correctness and completeness over a full day's data, not low latency. A nightly, or more frequent once the pipeline exists, batch export from the relational primary, or its read replica to avoid competing with production traffic, into a columnar warehouse is the standard shape for this.
Migration plan from a single-region relational database, staged to avoid a risky cutover.
- Stand up CDC off the existing single-region primary, with no application changes yet, feeding both a new per-region cache layer and a new analytics warehouse in parallel with whatever slower or absent read paths exist today, validated against real production traffic without being load-bearing.
- Cut checkout's availability-check reads over to the new regional cache behind a feature flag, region by region, while the commit and decrement still go through the unchanged relational primary. This isolates the riskiest part, the fast-read user experience path, from the correctness-critical commit, which stays untouched the longest.
- Cut analytics reporting over to the new warehouse, validate its output against the old direct queries on the primary during an overlap period, then retire the old direct-query path off production.
- Only after the above is stable, evaluate whether the commit step itself needs to become genuinely multi-region, a separate, higher-risk migration that should not be bundled into the read-path rollout above.
Worked example
flowchart LR
Checkout -->|fast local read| RegionCache[(Per-region cache: availability)]
Checkout -->|commit or decrement| RDBMS[(Relational primary: system of record)]
RDBMS -->|CDC| RegionCache
RDBMS -->|CDC, batch| Warehouse[(Analytics warehouse)]
Analyst[Replenishment analyst] -->|daily query| Warehouse
Trade-offs and pitfalls
A pitfall is treating the cache layer's staleness window as solved once it exists, without separately re-verifying the commit path still prevents overselling on its own. A fast, slightly stale read feeding a user-experience decision is fine; a fast, slightly stale read feeding the actual sale decision reintroduces the overselling bug in a new architecture.
A second pitfall is bundling the read-path migration and the write-path migration into a single cutover. They have very different risk profiles, worst case for the read path is a stale user experience, worst case for the write path is a financial or inventory correctness incident, which is exactly why the staged plan above deliberately separates them.
What would flip the recommendation: if checkout volume and user geography are actually concentrated in one region with only occasional international traffic, a simpler regional cache plus edge caching in front of the same single-region primary, without standing up a full multi-region write architecture, may hit the 50 millisecond target for the large majority of traffic at far lower operational cost. Size the actual geographic distribution before committing to full multi-region infrastructure.
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.