Cloud and Managed Database Services Questions
Running databases as managed cloud services rather than self-hosting them: choosing and provisioning managed relational and NoSQL offerings (for example Amazon RDS/Aurora/DynamoDB, Azure SQL Database/Cosmos DB, Google Cloud SQL/Spanner), sizing instances and storage, and designing an architecture (read replicas, connection proxies and pooling for managed or serverless compute) to meet a latency or availability target. It covers comparing provisioned versus serverless/autoscaling pricing and operational models, forecasting capacity and cost, and migrating a database onto or between managed deployments: moving a self-hosted database to a managed service, switching a workload from provisioned to serverless (or back), changing a single-AZ deployment to multi-AZ, and the questions to ask before adopting a new managed vendor. The core question is the split of responsibility and cost between the cloud provider and the engineering team: what the provider takes on (patching, infrastructure-level HA, automated backups) versus what remains the team's job (configuration, capacity and cost tuning, security posture like encryption and key rotation). It does not cover backup and restore mechanics, RPO/RTO planning, replication or sharding internals, or SQL-level query diagnostics.
A vendor offers a managed distributed SQL database that promises automatic sharding and rebalancing. List 6 operational questions you would ask before adopting it for a latency-sensitive transactional service.
Sample Answer
Direct answer
Six questions cover the ground a "handles sharding and rebalancing automatically" pitch conveniently skips: what latency the vendor will actually guarantee under your transaction shape, whether rebalancing degrades live traffic or pauses it, what the replication and consistency model implies for data loss and recovery time, how visible internal topology changes are to your own monitoring, how the system behaves when it cannot reach a quorum (a majority of replicas agreeing before a write counts as durable), and how hard it would be to leave. Ask all six framed against "latency-sensitive transactional," not as a generic vendor-diligence checklist: an answer that would satisfy an analytics workload can still be a bad fit here.
The six questions
| # | Question | Why it matters for a latency-sensitive transactional workload | Red flag in the vendor's answer |
|---|---|---|---|
| 1. Latency guarantee and measurement basis | "What 99th-percentile (P99) commit latency do you guarantee for a single-row transactional write, under a documented workload similar to mine, and is that backed by a service-level agreement (SLA) with a financial remedy?" | Distributed SQL databases add a consensus round trip (a network exchange where a quorum of replicas must acknowledge a write) on every commit. The number that matters is the tail latency of a small, single-row transaction, not an average and not a number measured on a large analytical batch. | Vendor quotes only mean or median latency, or a benchmark run on multi-row batch writes instead of the small single-row commits typical of a transactional service. |
| 2. Rebalancing impact on live traffic | "When the cluster rebalances automatically (a node joins, a hot shard splits, a failed replica is repaired), does the affected data stay fully read and write available, or does the system pause or throttle traffic on the piece being moved? Can rebalancing be rate-limited or scheduled away from peak hours?" | Rebalancing moves a shard (a contiguous chunk of the table's rows, grouped by key range) from one node to another by streaming a full copy of it across the network. Left unthrottled, that transfer competes with production traffic for network and CPU capacity right when the cluster's internal state is least predictable. This is not hypothetical: even mature distributed SQL engines have shipped rebalancing paths where an ordinary configuration change measurably degraded foreground query latency before the vendor patched it. | Vendor cannot describe a concrete rate-limiting or scheduling mechanism, or only offers "it is usually fine in practice." |
| 3. Replication, consistency, and disaster-recovery implications | "What is the default replication factor and quorum size per shard, what consistency guarantee does a client get by default (strict ordering across all replicas versus eventual convergence), and what recovery point objective (RPO, how much recent data could be lost) and recovery time objective (RTO, how long recovery takes) does that imply if a region fails mid-transaction?" | "Automatic sharding" says nothing about whether each shard's replicas live in one region (fast, but a regional outage can lose recent writes) or are kept in sync across regions by a global consensus protocol (safer, but adds real latency to every commit). | Vendor cannot state RPO or RTO as numbers, or uses "highly available" and "zero data loss" interchangeably. |
| 4. Observability into internal topology changes | "What real-time visibility exists into shard splits, leader elections (which replica currently owns the right to accept writes for a shard), and in-progress rebalancing, and can a latency spike in my own application performance monitoring (APM) tooling be correlated back to a specific internal event?" | If the database's internal reshuffling is invisible, a 2 a.m. latency spike looks identical to an application bug, and an on-call engineer burns an hour ruling out their own code before learning the database moved a shard underneath them. | Vendor offers only an aggregate cluster-health dashboard with no exportable, timestamped events for splits, elections, or rebalances. |
| 5. Failure-mode behavior | "During a network partition or a quorum loss on a shard, does the system fail closed (refuse writes, to stay consistent) or fail open (keep accepting writes on a stale replica)? Does failover for a lost leader happen automatically, or does it require a manual runbook step?" | For a transactional service, a database that fails open during a partition can accept writes that conflict or silently overwrite each other, surfacing as corrupted business state days later. That is worse than a short, honest availability gap. | Vendor cannot describe the failure mode precisely, or the answer implies writes are accepted during a partition "to keep availability high." |
| 6. Data portability and exit strategy | "If the contract ends, is export a native logical dump, a wire-protocol-compatible export another engine can read, or a proprietary snapshot only the vendor's own restore tooling understands? Who controls the encryption keys, and how long would extracting the full current data volume take at the vendor's stated export throughput?" | Automatic sharding and rebalancing are exactly the features that make a vendor's on-disk storage format the least portable. Verifying the exit path before signing is far cheaper than discovering it during a forced migration. | Vendor's export path is undocumented, or only supports restoring into their own product. |
Worked example
Take a payments authorization service: 3,000 transactions per second, a 40 ms P99 commit latency SLA of its own, and a requirement to survive a full region outage, so the team is evaluating a three-region deployment spanning, say, New York and San Francisco.
Question 1 asks the vendor for a P99 commit latency guarantee. Before even hearing their number, it is worth computing the physics floor a synchronous cross-region quorum write cannot beat, because no amount of engineering can get a commit acknowledgment back faster than light travels through fiber:
# Physics floor on cross-region consensus commit latency: a sanity check on any
# vendor SLA number, not a benchmark. All inputs are stated explicitly below.
great_circle_km = 4130 # approx. New York <-> San Francisco great-circle distance
routing_factor = 1.5 # real fiber paths run longer than the great circle
fiber_path_km = great_circle_km * routing_factor
speed_of_light_vacuum_km_s = 299_792
refractive_index_fiber = 1.5 # typical single-mode fiber
speed_in_fiber_km_s = speed_of_light_vacuum_km_s / refractive_index_fiber
one_way_ms = (fiber_path_km / speed_in_fiber_km_s) * 1000
round_trip_ms = one_way_ms * 2
print(f"fiber path length: {great_circle_km} km x {routing_factor} = {fiber_path_km:.0f} km")
print(f"speed in fiber: {speed_of_light_vacuum_km_s} / {refractive_index_fiber} = {speed_in_fiber_km_s:,.0f} km/s")
print(f"one-way propagation time: {fiber_path_km:.0f} km / {speed_in_fiber_km_s:,.0f} km/s = {one_way_ms:.3f} ms")
print(f"round trip (quorum ack floor): {one_way_ms:.3f} ms x 2 = {round_trip_ms:.2f} ms")
Output:
fiber path length: 4130 km x 1.5 = 6195 km
speed in fiber: 299792 / 1.5 = 199,861 km/s
one-way propagation time: 6195 km / 199,861 km/s = 30.996 ms
round trip (quorum ack floor): 30.996 ms x 2 = 61.99 ms
A round trip between New York and San Francisco alone costs about 62 ms, before any disk write, queueing, or server-side processing. If the vendor's answer to question 1 is "P99 commit latency under 40 ms" and question 3 reveals that every commit needs a quorum acknowledgment from a replica in the other region, those two answers contradict each other: the deployment cannot meet its own stated SLA. The honest fix is either to keep the write quorum inside one metro region (propagation delay measured in fractions of a millisecond) and replicate to the far region asynchronously for disaster recovery only, accepting a small, named RPO, or to relax the latency SLA. Either is a legitimate design choice; a vendor who does not surface this tension when questions 1 and 3 are asked together has not actually answered the question.
Trade-offs and pitfalls
The most common wrong turn is accepting a vendor's self-reported questionnaire as verification. A "yes" to all six questions from a sales engineer is not evidence: insist on a proof-of-concept load test that replays the actual read/write mix and transaction size against the vendor's cluster, sized to the real data volume, while watching P99 latency through a real rebalance (add a node, kill a replica, and observe). "Distributed" and "managed" are not synonyms for "safe for latency-critical transactions": consensus overhead, rebalancing behavior, and failure-mode defaults vary enormously between vendors, and that is exactly the corner of the evaluation a generic due-diligence checklist glosses over.
You want to reduce monthly costs by moving from provisioned RDS instances to a serverless model (Aurora Serverless v2) or DynamoDB on-demand. Create a migration and validation plan that addresses connection pooling differences, expected cold-start behaviors, schema or data-model changes, pricing model validation, performance testing, and rollback strategy.
Sample Answer
"Aurora Serverless v2" and "DynamoDB on-demand" are not two flavors of the same migration: moving RDS provisioned to Aurora Serverless v2 is a capacity-mode change on the same relational engine and the same schema, while moving to DynamoDB on-demand is a full re-platform onto a different data model with a real data migration and an application rewrite. The plan below targets the common case, RDS to Aurora Serverless v2, since that is what a single migration-and-validation plan covering connection pooling (the reusable set of open database connections an application keeps on hand instead of opening a new one per query), cold starts, schema changes, pricing, performance, and rollback actually maps onto cleanly; the DynamoDB path is addressed separately at the end because pretending it fits the same six-step plan would be misleading.
Migration and validation plan: RDS provisioned to Aurora Serverless v2
1. Confirm engine compatibility first, as its own step. Aurora Serverless v2's supported ACU (Aurora capacity unit) range depends on the specific engine version; some ranges (for example 0 to 256 ACU with scale-to-zero) require a minimum engine version. If the current instance is on an older version, upgrade the engine on the existing provisioned instance and validate that in isolation before touching capacity mode, so an engine-version regression and a capacity-mode regression are never conflated in the same change.
2. Connection pooling differences. Aurora sets the maximum connection ceiling based on the cluster's configured maximum ACU, not its current scaled-down capacity, specifically so active connections are not dropped when the cluster scales down. The practical risk runs the opposite direction from what people expect: if the previous provisioned instance was a large class with a generous default max_connections, and the new cluster's maximum ACU is set conservatively, the new ceiling could be lower than before. Check the application's connection pool size against the new maximum-ACU-derived ceiling before cutover, and keep RDS Proxy or an existing pooler in front regardless, since it smooths connection churn during scaling events even though Aurora states scaling itself does not disrupt existing connections.
3. Cold-start behavior. Normal Serverless v2 scaling within the 0.5-and-up ACU range has no cold start: scaling happens on the same host in most cases and does not wait for a quiet point or drop connections. The real cold-start risk only appears if the cluster is configured with a minimum capacity of 0 ACU (the "scale to zero" auto-pause feature on supported engine versions): resuming from a fully paused cluster does introduce a real resume delay. For an always-on production workload, do not set minimum capacity to 0; set it high enough to keep the working set (the rows and indexes your queries actually touch regularly) resident in the buffer pool (the database engine's in-memory cache of recently used data pages), which avoids both the auto-pause cold start and unnecessary cache-cold query latency (a query running slow because the data it needs isn't cached in memory and has to be read from disk instead) after every scale-down.
4. Schema and data-model changes. None are required, this is the entire advantage of choosing Serverless v2 over DynamoDB for a cost-reduction goal: it is the same PostgreSQL or MySQL engine underneath. What does need review is any parameter-group tuning that assumed a fixed, large instance class (for example manually tuned buffer pool or work_mem settings); Aurora auto-adjusts some capacity-linked parameters as the cluster scales, but custom overrides need to be re-validated against the new capacity range, not assumed to still be correct.
5. Pricing model validation. Do not trust a back-of-envelope estimate. Run the actual workload against a staging Serverless v2 cluster for at least one full representative traffic cycle (a week, to capture weekly cyclicality, longer if there's a known monthly peak) and compare the real ACU-hours consumed, at the current published Serverless v2 rate, against the current provisioned instance's hourly cost. The crossover point depends heavily on the workload's peak-to-average ratio: a low ratio (steady, always-busy traffic) usually still favors provisioned; a high ratio (bursty, with real idle periods) is where Serverless v2 wins. Validate with measured ACU utilization from staging, not a spec-sheet comparison.
6. Performance testing. Replay production-representative load, or better, a captured production query log, against the staging Serverless v2 cluster and compare p50/p95/p99 latency (the response time for the median, the slowest 5%, and the slowest 1% of requests, a way of measuring typical AND worst-case performance rather than just the average) against the current provisioned baseline. Specifically test the scaling boundary: force a load burst that crosses several ACU steps and confirm latency does not regress noticeably during the scale-up, since that is the one behavior genuinely new to this configuration.
7. Rollback strategy. Because the schema is unchanged, rollback is comparatively simple. Keep the original provisioned instance available (either still running in parallel during a defined validation window, or ready to be created back from the same cluster since Aurora supports converting a Serverless v2 writer back to a fixed provisioned instance class), and define the rollback trigger criteria up front, specific latency or cost thresholds, rather than deciding what counts as "bad enough to roll back" during an actual incident.
If the real target is DynamoDB on-demand instead
Treat this as a separate, larger project, not a variant of the plan above. It requires a genuine data-model redesign (partition and sort keys replacing joins, foreign keys, and ad hoc filtering), a real data migration (for example via AWS Database Migration Service or a custom extract-transform-load pipeline, not a capacity-mode toggle), and rewriting every query path that relied on relational joins. Two of the six items above barely apply: DynamoDB uses request-based SDK calls over HTTP rather than persistent connections, so "connection pooling" in the RDS sense mostly does not exist, and there is no cold-start concept analogous to Aurora Serverless v2's auto-pause. Forcing this migration into the same six-step template as the Aurora path would hide the fact that it is a fundamentally different, higher-effort, higher-risk change.
Trade-offs and pitfalls
- Cost reduction is not the same as risk reduction. The Aurora Serverless v2 path is low schema risk but carries real behavioral risk: a misconfigured connection ceiling or a minimum ACU set too low for the working set can both show up as production latency regressions even though "nothing about the schema changed."
- Bundling the engine-version upgrade and the capacity-mode switch into one change removes the ability to tell which change caused a regression if one appears; do them as two separate, separately validated steps.
- Picking DynamoDB "because it's serverless and cheaper" without first validating that the application's access patterns actually fit a key-value model can produce an outcome worse than simply staying on provisioned RDS: a relational workload force-fit onto DynamoDB tends to reappear as either a pile of expensive scans or an application layer quietly reimplementing joins in code.
Evaluate the trade-offs of using a managed multi-region database service (for example Cloud Spanner, Cosmos DB, or Aurora Global Database) versus running a self-managed sharded cluster on VMs or Kubernetes, for a system handling roughly 10 million QPS reads and 200,000 writes per second. Cover latency, consistency, operational complexity, cost, scaling behavior, and disaster-recovery capabilities, and give your recommendation.
Sample Answer
Direct answer
At 10 million reads per second and 200,000 writes per second, the default recommendation is a managed multi-region service (Spanner-class horizontal partitioning or Cosmos DB-class tunable consistency) over a self-managed sharded cluster, unless the team already operates comparable distributed stateful infrastructure in-house and the consistency requirements are loose enough to accept asynchronous cross-region replication. At this scale the operational surface a managed vendor absorbs, quorum health (whether enough replicas are up and agreeing to keep accepting writes), rebalancing, and cross-region failover (automatically switching traffic to a healthy region if one goes down) across what turns out to be hundreds of nodes, is exactly the class of problem these services are purpose-built for. The worked example and framework below show why, and name what would flip the recommendation.
Framework: six axes
| Axis | Managed multi-region (for example Cloud Spanner, Cosmos DB, Aurora Global Database) | Self-managed sharded cluster (VMs or Kubernetes) |
|---|---|---|
| Latency | In-region reads and writes stay low: Cosmos DB's own documentation guarantees a 99th-percentile (P99) latency under 10 ms for reads and writes within a region, at any consistency level. Cross-region strong consistency is bounded by physics, not engineering: Cosmos DB states cross-region strong-consistency write latency as two times the round-trip time (RTT) between the two farthest regions, plus 10 ms P99. | Fully tunable: keep the write quorum inside one metro region for low latency, and replicate elsewhere asynchronously for disaster recovery (DR) only, if the business can accept that. Requires the team to design and prove this topology correct itself. |
| Consistency | Comes largely built in. Cosmos DB offers five levels (Strong, Bounded Staleness, Session, Consistent Prefix, Eventual). Spanner gives a single strong guarantee, external consistency (global linearizability, meaning all observers agree on the order transactions actually committed in), via TrueTime, a globally synchronized clock with a bounded uncertainty window that every commit waits out before acknowledging. | Owned entirely by the team. Single-shard transactions are ordinary single-node ACID (atomicity, consistency, isolation, durability). Cross-shard transactions need a distributed-transaction layer, for example two-phase commit (a protocol that coordinates an atomic commit across multiple shards), that the team designs, implements, and must prove correct under partition and retry. |
| Operational complexity | Patching, quorum membership, internal rebalancing, and backup mechanics are the vendor's job. The team owns schema design, indexing, provisioned or consumed throughput sizing, and query-level tuning. | Everything: failover automation, connection pooling and proxy layers, monitoring, and on-call runbooks for partition and quorum-loss scenarios are the team's to write and operate, with no vendor site reliability engineering (SRE) team absorbing pages. |
| Cost | Priced at a premium over commodity compute (consumption-based request units, or per-node or per-vCPU list pricing depending on the specific service and tier; check the current pricing page for the exact service and region before budgeting, since list prices change without a fixed schedule). The premium buys engineering time the team does not have to spend building the same capability. | Commodity VM or Kubernetes compute, storage, and network egress rates, but the true cost includes the engineering headcount to build and operate the sharding, failover, and DR tooling a managed service ships by default. That cost does not appear on the cloud bill. |
| Scaling behavior | Elastic and mostly automatic. Aurora Global Database adds up to 10 read-only secondary regions with storage-based replication typically under a second of lag, though only the primary region accepts direct writes (a "write forwarding" feature lets secondaries relay writes to the primary, at the cost of that round trip). Spanner and Cosmos DB partition writes horizontally across internal shards transparently to the application. | Fully within the team's control (shard key choice, resharding strategy, instance sizing), but every reshard is a project the team runs and risks itself, exactly the work automatic rebalancing would otherwise absorb. |
| Disaster recovery | Documented, tested guarantees. A Cloud Spanner multi-region configuration places two read-write regions (two replicas each) plus a witness region (a fifth voting replica that participates in the write quorum but does not serve reads), five voting replicas total; a write quorum needs the leader region's replica plus any two of the other four, so the configuration survives the complete loss of any one region with zero data loss, a recovery point objective (RPO) of zero. Aurora Global Database offers a planned switchover (no data loss) and an unplanned failover after a regional outage (a small, non-zero RPO bounded by the sub-second replication lag (the delay between a write completing in the primary region and arriving in the secondary) at the moment of failure). | Fully bespoke: cross-region async replication, backup and restore, or a warm standby are all buildable, but the recovery time objective (RTO) and RPO the system actually delivers are only as good as the team's own runbook and the disaster tests actually run, not a documented product guarantee. |
Worked example
Whichever path is chosen, the raw infrastructure needed at this scale is a useful reality check: 10 million reads/sec and 200,000 writes/sec is real online transaction processing (OLTP) hardware regardless of who operates it.
import math
# Inputs: the QPS figures come from the question; the rest are stated,
# illustrative capacity-planning assumptions (NOT vendor-measured benchmarks).
read_qps = 10_000_000
write_qps = 200_000
read_capacity_per_node = 50_000 # assumed sustained cached-read throughput of one well-provisioned OLTP node
write_capacity_per_shard = 8_000 # assumed sustained write throughput of one shard's single-writer primary
dr_regions = 3 # regions in the deployment
shards_needed = math.ceil(write_qps / write_capacity_per_shard)
total_read_nodes = math.ceil(read_qps / read_capacity_per_node)
read_nodes_per_shard = math.ceil(total_read_nodes / shards_needed)
nodes_per_shard_group = 1 + read_nodes_per_shard # 1 primary + its read nodes
total_nodes_single_region = shards_needed * nodes_per_shard_group
total_nodes_all_regions = total_nodes_single_region * dr_regions
print(f"shards needed (write-bound): {write_qps:,} / {write_capacity_per_shard:,} -> {shards_needed}")
print(f"read-serving nodes needed (read-bound): {read_qps:,} / {read_capacity_per_node:,} -> {total_read_nodes}")
print(f"read nodes per shard: {total_read_nodes} / {shards_needed} -> {read_nodes_per_shard}")
print(f"nodes per shard group (1 primary + reads): {nodes_per_shard_group}")
print(f"total nodes, single region: {shards_needed} x {nodes_per_shard_group} -> {total_nodes_single_region}")
print(f"total nodes, {dr_regions}-region deployment: {total_nodes_single_region} x {dr_regions} -> {total_nodes_all_regions}")
Output:
shards needed (write-bound): 200,000 / 8,000 -> 25
read-serving nodes needed (read-bound): 10,000,000 / 50,000 -> 200
read nodes per shard: 200 / 25 -> 8
nodes per shard group (1 primary + reads): 9
total nodes, single region: 25 x 9 -> 225
total nodes, 3-region deployment: 225 x 3 -> 675
Roughly 675 stateful nodes across three regions, on these stated assumptions. That number does not change based on who manages it: a managed multi-region service still provisions and rebalances something in that range internally, it simply does not page the team when a node in it fails. The real question is whether the team wants to own the failure domain for 675 stateful nodes directly, given that a self-managed cluster at this size effectively requires building an internal SRE function for exactly that fleet.
flowchart LR
subgraph MG["Managed multi-region service"]
C1["Client"] --> R1["Regional endpoint"]
R1 --> P1[("Primary region")]
P1 -. replicates .-> S1[("Secondary region A")]
P1 -. replicates .-> S2[("Secondary region B")]
end
subgraph SM["Self-managed sharded cluster"]
C2["Client"] --> RT["Shard router"]
RT --> SH1[("Shard 1 primary")]
RT --> SH2[("Shard 2 primary")]
SH1 --> RE1[("Shard 1 replica")]
SH2 --> RE2[("Shard 2 replica")]
end
In the managed topology, the vendor owns everything to the right of the regional endpoint: which replica leads, how replication and quorum work, how a lost region is detected and routed around. In the self-managed topology, the team owns and must operate the shard router, the primary and replica wiring per shard, and the logic for what happens when a shard's primary disappears.
Recommendation and what would flip it
Recommend the managed multi-region path by default at this scale. The deciding factors are the team's existing distributed-systems operational maturity (does it already run something like Vitess- or CockroachDB-class infrastructure at hundreds of nodes today), the actual consistency requirement (does the product need synchronous cross-region strong consistency, or would single-writer-per-shard with asynchronous DR replication genuinely be acceptable), and whether the cost delta at roughly 675 nodes' worth of infrastructure materially matters to the budget, since the managed premium multiplies across a large fleet at this scale rather than a handful of instances.
It flips to self-managed when all three line up: the team already operates comparable stateful infrastructure in-house, the workload tolerates a single-writer-per-shard model with asynchronous cross-region replication instead of synchronous strong consistency, and the cost savings at this node count are large enough to justify paying the engineering-time cost once, building the sharding, failover, and DR tooling, instead of paying the managed-service premium every month indefinitely. Defaulting to "managed is always safer" ignores that at 675 nodes the managed premium is not a rounding error, and a team that already has the operational muscle for a fleet this size may reasonably prefer to keep that premium as engineering budget instead.
Trade-offs and pitfalls
The dominant failure mode on this kind of question is picking a side without naming the deciding variable, or hedging into "it depends" without ever committing. A second common mistake is assuming managed multi-region services give multi-region strong consistency for free: Cosmos DB explicitly does not support strong consistency on a multi-region-write account, and even Spanner's external consistency costs real latency, a bounded commit-wait on every write, it is not free ordering. A third is treating "self-managed" as automatically cheaper: at 675 nodes, the engineering cost of building and operating correct distributed-transaction and failover logic in-house is usually larger, not smaller, than the managed premium, unless that capability already exists in the organization.
You're designing the managed-database architecture for an ad-serving platform with 100M monthly active users. Characteristics: 95% reads, 5% writes; p95 read latency target < 50ms globally; peak reads 50k RPS; writes are small updates ~10k TPS during peak. Data shape: user preferences and campaign state, schema somewhat flexible. Recommend a managed database architecture (products and components) to meet the latency SLA and discuss caching, indexing, and operational concerns.
Sample Answer
This is a key-value access pattern (look up preferences and campaign state by user or campaign ID) at extreme read scale with a hard global latency target, so the recommendation is DynamoDB as the system of record, fronted by DynamoDB Accelerator (DAX, an in-memory read cache purpose-built for DynamoDB) in each region, replicated across regions with DynamoDB global tables, behind latency-based routing that sends each user to their nearest region.
flowchart LR
Client[Ad-serving client] --> Router[Latency-based routing]
Router --> AppUS[App tier, us-east-1]
Router --> AppEU[App tier, eu-west-1]
AppUS -->|read, cache-first| DAXUS[DAX cluster, us-east-1]
AppEU -->|read, cache-first| DAXEU[DAX cluster, eu-west-1]
DAXUS -->|cache miss| DDBUS[(DynamoDB table, us-east-1)]
DAXEU -->|cache miss| DDBEU[(DynamoDB table, eu-west-1)]
DDBUS <-.async global replication.-> DDBEU
Why this architecture fits the workload
- 95% reads, flexible schema: DynamoDB's item model (no fixed columns per row) fits user preference and campaign-state documents that vary in shape without requiring a schema migration every time a new attribute is added, and it is purpose-built for very high, uniformly-distributed read throughput.
- p95 read latency (the response time that 95% of requests come in at or under; the slowest 5% are allowed to exceed it) under 50ms globally: a single-region database cannot hit this for users far from it; global tables replicate the same table to multiple regions so every region has a local, complete copy, and DAX gives single-digit-millisecond to microsecond reads for cached keys on top of that.
- Writes are small, targeted updates: DynamoDB's per-item write model (not full-row rewrites, not multi-table joins) matches "update this user's preferences" or "update this campaign's remaining budget" well, and 5% writes at 10k TPS peak is comfortably within what a well-partitioned table handles.
Capacity math for the stated peak load
Assume a 1 KB average item size (a compact preference or campaign-state document) and a 90% DAX cache hit ratio for reads, both reasonable for a keyed lookup workload with a lot of repeat access to the same active users and campaigns. DynamoDB prices reads in read capacity units (RCU): a strongly consistent read of an item up to 4KB costs 1 RCU, but an eventually consistent read of the same item costs half that, 0.5 RCU, because DynamoDB can answer it from any replica without first confirming it has the very latest write. A cached preference or campaign-state read here does not need to reflect a write from the last few hundred milliseconds, so eventually consistent reads are the right, cheaper default; strongly consistent reads would double the RCU figures below:
peak_read_rps = 50_000
peak_write_tps = 10_000
dax_cache_hit_ratio = 0.90
RCU_PER_PARTITION = 3000 # AWS-documented per-partition ceiling (read units/sec)
WCU_PER_PARTITION = 1000 # AWS-documented per-partition ceiling (write units/sec)
rcu_per_read = 0.5 # eventually consistent read, item fits in one 4KB chunk
wcu_per_write = 1 # 1KB item fits in one 1KB chunk
backend_reads_per_sec = peak_read_rps * (1 - dax_cache_hit_ratio)
backend_rcu_per_sec = backend_reads_per_sec * rcu_per_read
worst_case_rcu_per_sec = peak_read_rps * rcu_per_read # DAX cold, e.g. right after a deploy
write_wcu_per_sec = peak_write_tps * wcu_per_write
print(f"Backend reads/sec with warm DAX: {backend_reads_per_sec:,.0f}")
print(f"Backend RCU/sec with warm DAX: {backend_rcu_per_sec:,.0f}")
print(f"RCU/sec if DAX is cold: {worst_case_rcu_per_sec:,.0f}")
print(f"WCU/sec (writes, uncacheable): {write_wcu_per_sec:,.0f}")
print(f"Min partitions from writes: {-(-write_wcu_per_sec // WCU_PER_PARTITION):.0f}")
print(f"Min partitions from cold-cache reads: {-(-worst_case_rcu_per_sec // RCU_PER_PARTITION):.0f}")
Backend reads/sec with warm DAX: 5,000
Backend RCU/sec with warm DAX: 2,500
RCU/sec if DAX is cold: 25,000
WCU/sec (writes, uncacheable): 10,000
Min partitions from writes: 10
Min partitions from cold-cache reads: 9
With a warm DAX cache, the DynamoDB table itself only has to absorb roughly 2,500 RCU/sec (well under one partition's ceiling), meaning the table's real sizing constraint here is the write path, which needs throughput spread across at least 10 partitions (10,000 WCU/sec divided by the documented 1,000 WCU-per-partition ceiling). On-demand capacity mode is the right choice at this scale rather than manually provisioned throughput, since the traffic pattern in an ad-serving platform (campaign launches, dayparting) is not perfectly steady, and on-demand removes the need to hand-tune provisioned throughput and auto scaling targets across two regions.
Caching, indexing, and operational concerns
- Partition key design: key user preference items by a well-distributed attribute like
user_id(100M distinct values spreads load evenly). Do not key campaign-state items bycampaign_idalone if a small number of high-spend campaigns dominate traffic; that concentrates writes on a few partitions and creates a hot partition even though the table's aggregate throughput looks fine. Use write sharding (a suffixed key likecampaign_id#shard_n) for any campaign whose write rate alone could approach the 1,000 WCU-per-partition ceiling. - DAX: run a DAX cluster with at least 3 nodes across availability zones (AZs) per region for high availability, and set a write-through or short time-to-live (TTL, how long a cached value is trusted before being refetched) policy on campaign-state items specifically, since a stale campaign budget or targeting rule served from cache for too long has real business cost (overspend, mistargeting), unlike a slightly stale user preference.
- Global tables and consistency: replication between regions is asynchronous and uses last-writer-wins conflict resolution, so a user who is rerouted between regions (failover, or a mobile client crossing a region boundary) can briefly see a slightly older version of their own last write. This is acceptable for preferences and most campaign state, but if a specific field (for example remaining campaign budget) must never double-spend across regions, that field needs a different pattern, for example routing all writes for a given campaign to a single "home" region rather than treating every region as equally writable.
- Monitoring: alarm on
ThrottledRequestsandConsumedWriteCapacityUnitsper table (and per partition-level hot-key indicators), DAX cache hit rate (a real hit ratio far below the assumed 90% means the 2,500 RCU/sec sizing estimate above is wrong and the table needs to absorb more direct traffic), and cross-region replication latency for global tables. - Indexing beyond the primary key: add global secondary indexes (GSIs) only for query patterns that genuinely need them (for example looking up all campaigns for an advertiser), since every GSI is itself a separately-partitioned, separately-throughput-provisioned structure with its own hot-key risk, not a free index the way a relational database's secondary index is.
Explain connection pooling and why it matters for managed databases, especially with serverless compute (AWS Lambda) and managed MySQL/Postgres. Name pooling solutions you'd consider, explain how you'd size the pool, and describe what you'd configure to avoid connection storms.
Sample Answer
Connection pooling reuses a small set of already-open database connections across many requests instead of opening a fresh connection (a TCP handshake plus authentication) for every request. It matters most acutely with serverless compute like AWS Lambda because each concurrent Lambda execution environment behaves like an independent short-lived process: if your code opens a database connection on invocation, a burst of concurrent invocations opens that many connections nearly simultaneously, and a managed MySQL or PostgreSQL instance has a hard max_connections ceiling tied to its instance size. A traffic spike that is entirely reasonable in query volume can still take the database down purely on connection count.
Why it matters beyond Lambda
Every open connection costs the database real memory and, for PostgreSQL specifically, a whole backend OS process, so connections are expensive per-unit even when idle. Opening and closing them constantly (rather than reusing a pool) adds latency and CPU overhead from repeated authentication and TLS negotiation on every request, on top of the risk of simply running out of connection slots.
Pooling solutions to consider
- RDS Proxy: AWS-managed, sits between the application and RDS (Amazon's managed relational database service for engines like MySQL and PostgreSQL) or Aurora (AWS's own MySQL- and PostgreSQL-compatible managed database with a different underlying storage engine). It pools and multiplexes connections, supports IAM authentication and credentials via AWS Secrets Manager, and automatically fails over to a standby while preserving application connections. It can only be attached to a writer instance, not a read replica, and it requires the proxy and the database to be in the same virtual private cloud (VPC).
- PgBouncer: self-run, the standard external pooler for PostgreSQL, used constantly in front of Lambda-to-Postgres workloads because of its transaction-mode pooling.
- ProxySQL: the MySQL equivalent, adding query routing and caching on top of pooling.
- Cloud SQL Auth Proxy / connectors: Google Cloud's equivalent pattern for Cloud SQL, handling authenticated, pooled connections from serverless callers.
Sizing the pool
Size the pool from actual concurrency, not guesswork, using Little's Law: the average number of concurrent database connections a workload needs equals its request rate multiplied by the average time each request holds a connection.
L=λ×Wrequests_per_sec = 2000 # peak invocation rate hitting the database
avg_query_duration_s = 0.015 # 15ms average time from "acquire connection" to "commit/release"
concurrent_db_connections_needed = requests_per_sec * avg_query_duration_s
db_max_connections = 1000 # example ceiling for a large Aurora instance class
usable_max_connections = db_max_connections - 50 # reserve headroom for admin/replication
headroom_multiplier = 3 # slack for variance and short bursts above the mean
recommended_pool_size = concurrent_db_connections_needed * headroom_multiplier
print(f"Concurrent DB connections needed on average: {concurrent_db_connections_needed:.1f}")
print(f"Recommended pool size with {headroom_multiplier}x headroom: {recommended_pool_size:.0f}")
print(f"That is {recommended_pool_size/usable_max_connections:.1%} of usable max_connections")
Concurrent DB connections needed on average: 30.0
Recommended pool size with 3x headroom: 90
That is 9.5% of usable max_connections
The point this makes concrete: even at 2,000 requests/sec, if each request only holds a connection for 15ms, the database only ever needs about 30 connections open at once on average, 90 with generous headroom. A pooler lets a huge number of Lambda-side clients share that small, stable set of database-side connections, which is the entire reason pooling and Lambda go together.
What to configure to avoid connection storms
- Use transaction-mode pooling, not session mode. In transaction mode, the database-side connection is returned to the pool the instant a transaction commits, not held for the life of the client's connection to the pooler. This is what lets 30 to 90 real database connections serve thousands of Lambda invocations: each invocation borrows a connection for milliseconds, not for its entire lifetime.
- Point Lambda only at the pooler's endpoint, never at the database engine directly. If any code path bypasses the pooler, it reintroduces the exact one-connection-per-invocation problem pooling exists to solve.
- Set the pooler's client-facing limit high but its database-facing limit bounded near the database's real capacity: thousands of client-side slots, tens to low hundreds of database-side ones.
- Set idle and statement timeouts so an abandoned or slow Lambda-side connection cannot hold a pooled database slot indefinitely.
- Avoid patterns that force "pinning." RDS Proxy (and poolers generally) fall back to a dedicated, unshared connection for a session that changes connection-level state, for example a session-level
SETstatement, an unclosed transaction, or, on RDS Proxy specifically, any statement larger than 16 KB. A pinned session stops multiplexing entirely for its lifetime, which quietly reintroduces the connection-storm risk the pooler was meant to prevent.
Trade-offs and pitfalls
- A pooler is another moving part with its own failure modes: if it goes down, so does every connection through it, so it needs to be highly available in its own right (RDS Proxy handles this for you; a self-run PgBouncer needs its own redundancy plan).
- Transaction-mode pooling breaks session-level features that assume a stable connection: prepared statements that persist across queries, session-level temp tables, and
LISTEN/NOTIFYin Postgres either don't work or behave differently. Audit the application for these before switching modes. - Pooling fixes connection exhaustion; it does not fix a database that is genuinely CPU- or IO-bound. If queries are slow, a pool just queues requests instead of rejecting them outright, which can turn a fast failure into a slow, harder-to-diagnose one if queue depth isn't monitored.
Unlock Full Question Bank
Get access to all 20 Cloud and Managed Database Services interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.