Observability and Monitoring Architecture Questions
Building visibility into infrastructure and services: metrics, logs, and traces, dashboards and alerting, SLIs/SLOs, and the design of an observability stack. Covers instrumenting systems for actionable signal, reducing alert noise, and diagnosing production issues from telemetry. Infrastructure-wide observability, distinct from network-specific monitoring.
You're designing monitoring for a Kubernetes platform that mixes stateless front-ends with stateful databases. Decide which components should run as DaemonSets, which as sidecars, and which as centralized services, and explain how you'd minimize resource overhead on the nodes running the stateful workloads without losing signal fidelity.
Sample Answer
Direct Answer
Put uniform, low-cost collection (node and infra metrics, general logs) on a DaemonSet everywhere, including on the stateful nodes, since a flat per-node cost is easy to reason about and budget for. Reserve sidecars for signal that genuinely requires in-process, per-pod access a node agent cannot get (an application's internal connection pool state, for example), and be deliberate about NOT putting a sidecar on the database pods unless that condition is actually met. Anything needing cross-pod correlation or heavy processing (tail sampling, log parsing, deep query-stats collection) belongs off-node entirely, in a centralized service that polls or receives from the stateful workload remotely, so it never competes with the database for the node's own CPU and memory. This reasoning holds the same way on a managed control plane (EKS, GKE, AKS) as on self-managed Kubernetes, since it is about resource contention on the node, not about who runs the API server.
Structured Elaboration
Placement topology
flowchart LR
subgraph STATELESS["Stateless Front-End Nodes"]
FE["Front-end Pods"] --> DS1["DaemonSet Agent"]
end
subgraph STATEFUL["Stateful DB Nodes, capacity-constrained"]
DB[("Database Pod")]
DS2["DaemonSet Agent: node metrics only"]
end
DS1 --> GATEWAY["Central Collector Gateway"]
DS2 --> GATEWAY
POLLER["Off-node DB Metrics Poller"] -->|"remote scrape"| DB
POLLER --> GATEWAY
GATEWAY --> BACKEND[("Backend")]
Stateless front-end nodes
Lower stakes: DaemonSet for infra metrics and logs, sidecars are affordable here if a particular front-end service wants per-pod tracing detail, since these nodes typically run with more spare headroom.
Stateful database nodes
Higher stakes: these nodes are usually provisioned close to capacity by the database's own resource requests, leaving little headroom for anything else. The DaemonSet's flat, small, predictable overhead is safe here. A sidecar doing anything beyond trivial passthrough is risky, because its resource usage now directly competes with the database's own background work (vacuum, checkpointing, WAL flush) for the same constrained node.
What to centralize instead of running on the stateful node
Heavy or periodic collection against the database (query-stats dumps, slow-query log parsing) should run as an off-node poller hitting the database's metrics endpoint remotely. This adds a network hop and a small amount of latency to seeing that data, but adds zero incremental CPU or memory pressure on the node itself beyond the lightweight DaemonSet.
Worked Example
Assume a stateful node with 16 vCPU total, where the database pod's own resource request is 14 vCPU, leaving 2 vCPU (2,000m) of headroom for everything else on the node.
The baseline DaemonSet agent (node and infra metrics only) requests 200m CPU:
2,000m200m=10% of remaining headroomThat is a safe, predictable cost. Now add a hypothetical sidecar that periodically runs a heavier in-process collection task, bursting to 500m CPU during its collection window:
2,000m200m+500m=2,000m700m=35% of remaining headroomA 35% claim on the database's only spare capacity, right when the sidecar's own collection burst happens to run, is exactly the kind of contention that can visibly delay the database's own background maintenance work. Moving that heavier collection off-node (a remote poller scraping the database's metrics endpoint) removes the 500m burst from the node entirely, leaving the node's overhead at the flat 10% DaemonSet baseline regardless of collection frequency.
Trade-offs and Pitfalls
The off-node poller trades node-local safety for a network dependency and slightly staler data (whatever the poll interval is, versus in-process real-time access). For a database, this is almost always the right trade: a few seconds of staleness on query-stats data is a much smaller risk than periodic CPU contention with the database engine itself.
Not every stateful workload has this much headroom to spare. If a database's own resource request already consumes 15.5 of 16 vCPU, even the flat DaemonSet overhead becomes a real constraint, and the answer shifts toward running collection on a dedicated node pool or accepting sparser, lower-frequency node-level metrics on those specific hosts rather than the fleet-wide default cadence.
A common mistake is applying the same sidecar-friendly policy used for stateless front-ends uniformly across the whole cluster, without re-evaluating it against the stateful nodes' actual spare capacity. The placement decision should be driven by measured headroom per node pool, not by a single fleet-wide policy applied blindly.
Design a metrics ingestion pipeline that must accept roughly one million data points per second across three regions. Cover collector and agent placement, buffering and batching, message broker selection and partitioning keys, deduplication, backpressure handling, fault tolerance, and where you would perform pre-aggregation or rollups to reduce load downstream.
Sample Answer
Direct Answer
Split ingestion by region so no single path crosses a WAN in the hot loop: collectors batch and buffer locally, hand off to a partitioned durable log keyed by series identity, and a stream layer deduplicates and rolls up before anything touches long-term storage. The three levers that make one million points per second tractable are partition count (parallelism), batch size (write amplification), and pre-aggregation (what actually needs to survive at full resolution).
Structured Elaboration
Pipeline topology
flowchart LR
subgraph REGION["Per-Region Tier (x3)"]
APP[Service Instances] --> AGENT[Collector Agent]
AGENT --> WAL[("Local WAL Buffer")]
WAL --> KAFKA[["Kafka: 32 partitions by tenant+series key"]]
end
KAFKA --> SP["Stream Processor: dedup + rollup"]
SP --> TSDB[("Hot TSDB")]
SP --> OBJ[("Object Storage: rollups")]
Collector and agent placement
Run a lightweight collection tier per region (behind a regional load balancer, autoscaled in Kubernetes) so no metric point leaves its region before being durably buffered. Cross-region replication happens downstream, at the storage layer, never in the write-critical path, so a WAN blip in one region does not add latency to the other two.
Buffering and batching
Each collector holds a local disk-backed buffer (a WAL, RocksDB-backed queue, or equivalent) so a restart or a downstream stall does not drop in-flight data. Batch by size, not purely by time: a size trigger keeps latency proportional to actual load instead of always waiting out a fixed window.
Message broker and partitioning keys
Use a partitioned durable log (Kafka or equivalent) per region. Partition key: hash of tenant_id + metric_name. This keeps every point for a given series in the same partition, which is what makes per-series ordering and local windowed aggregation possible downstream, at the cost of potential hot partitions for very high-cardinality single tenants.
Deduplication
Give every point an idempotency key (source_id + monotonic sequence number). The stream processor keeps a rolling window of seen keys; anything already seen in that window is a duplicate produced by a retry, not new data.
Backpressure handling
The collector never blocks its callers. When the local buffer approaches capacity, it sheds lowest-priority series first (a stated priority tier, not silent random drop) and raises an explicit metric so the shedding is visible, not a silent gap in a dashboard three weeks later.
Fault tolerance
Replicate the log (replication factor 3) across brokers in multiple availability zones. Collectors are stateless except for the local buffer, so a lost collector instance loses only unflushed buffer content, not history. Stream processors checkpoint offsets so a crash resumes from the last committed point, not from zero.
Pre-aggregation and rollups
Do windowed rollups (count, sum, min, max, and a percentile sketch) in the stream layer before the write to the hot store. Keep raw resolution for a short window (operators debugging an active incident need seconds-level data); roll everything older than that window up to coarser resolution, since almost no dashboard or alert needs one-second granularity on data from an hour ago.
Worked Example
Assume 1,000,000 points/sec split across 3 regions and an average encoded point size of 150 bytes (timestamp, value, and label set after protobuf encoding, a stated design input, not a benchmark). Regional steady load is then:
rateregion=31,000,000≈333,333 pts/s⇒333,333×150 B≈50 MB/sPlan for regional skew: if one region fails, its traffic can fail over to the nearest healthy region, so size for up to 1.5x the average, 500,000 pts/s = 75 MB/s peak.
Partition count. Choose a conservative per-partition write budget of 10 MB/s (accounts for replication-factor-3 fsync overhead on commodity brokers, a stated assumption, not a vendor benchmark):
partitionssteady=⌈1050⌉=5,partitionspeak=⌈1075⌉=8Round up to 32 partitions per regional topic: this gives 4x headroom over the peak-throughput floor of 8, so future traffic growth or a temporary partition hot-spot does not require an emergency repartition, and it splits evenly across, say, 8 stream-processor instances at 4 partitions each.
Batch fill time. At the peak-derived floor of 8 partitions, each carries roughly 50 MB/s / 8 = 6.25 MB/s. A 1 MB size-triggered batch fills in:
tbatch=6.25 MB/s1 MB=0.16 s=160 msThat is well under a reasonable 500 ms linger cap, so the size trigger (not the timeout) governs flush cadence under normal load, and the timeout only matters for low-traffic partitions.
Deduplication memory. With a 2-minute (120 s) dedup window at the regional steady rate, the number of distinct keys the stream processor must track at once is:
n=333,333×120≈4×107 keysSizing a Bloom filter (a compact structure that answers "have I possibly seen this key before," with a small, tunable false-positive rate but no false negatives) for these keys at a 0.1% false-positive rate (p=0.001):
m=(ln2)2−nlnp≈0.48054×107×6.908≈5.75×108 bits≈71.9 MBA roughly 72 MB Bloom filter per region is cheap enough to keep fully in memory and gives a bounded, known false-positive rate for the dedup layer, versus an exact hash-set which would need far more memory to track 40 million live keys.
Local buffer disk sizing. To survive a 10-minute (600 s) broker outage at the regional steady rate without dropping data:
bufferdisk=50 MB/s×600 s=30,000 MB=30 GBProvisioning roughly 30 GB of local disk per regional collector tier is the concrete number that backs the "collectors survive a broker outage" claim, not just an assertion that buffering exists.
Trade-offs and Pitfalls
A common alternative is writing directly from collectors to the time-series store, skipping the durable log entirely. That removes a hop and its operational cost, but couples ingest availability directly to storage availability: any storage hiccup now blocks collectors instead of just delaying a downstream consumer. The durable log is worth its cost specifically because it decouples those failure domains.
Pre-aggregation trades ingest-time compute for downstream storage and query cost: computing rollups at 1,000,000 points/sec needs real CPU budget in the stream layer, and if the rollup logic has a bug, it is much harder to recompute correct history than if raw data were simply sitting untouched in cheap storage. Keep raw data for a bounded window specifically so a bad rollup is recoverable.
The partition key of tenant_id + metric_name preserves per-series locality but can create a hot partition if one tenant emits a disproportionate share of traffic; watch for this and consider adding a bucket suffix to the key for known outlier tenants rather than repartitioning the whole topic reactively.
Describe architectural patterns to make a telemetry ingestion pipeline resilient to backpressure from downstream storage, for example when the time-series database becomes temporarily unavailable or traffic spikes 10x during an incident. Cover buffering, rate-limiting, circuit breakers, retry strategy, and how you would surface the pipeline's own health to the teams depending on it.
Sample Answer
Direct Answer
Decouple the pipeline from the sink with a durable buffer so a slow or unavailable time-series database does not propagate latency back to producers, wrap writes to the sink in a circuit breaker so a struggling database is not also hammered by retries, and treat the pipeline's own saturation state (queue depth, drop rate, breaker state) as a metric other teams can see, not an internal detail that only shows up as "my dashboard is missing data" after the fact.
Structured Elaboration
Resilience pipeline
flowchart LR
PROD["Producers"] --> RL["Rate Limiter"]
RL --> BUF[("Durable Buffer")]
BUF --> CB{"Circuit Breaker"}
CB -->|"closed"| TSDB[("TSDB Writer")]
CB -->|"open"| RETRY["Backoff + Retry"]
RETRY -.-> CB
BUF --> HEALTH["Queue-Depth / Drop-Rate Metrics"]
HEALTH --> DASH["Status Dashboard"]
Buffering
A durable, partitioned queue (Kafka, or an equivalent persistent broker) sits between collection and the storage writer. Producers write to the queue and get an ack independent of whether the writer is keeping up, which is what actually decouples the two.
Rate-limiting
A token-bucket limiter in front of the writer caps how fast it attempts to push into storage, so a recovering database is not immediately re-flooded the moment it comes back up. The bucket's burst capacity should match what the buffer can absorb, not an arbitrary number.
Circuit breaker
Wrap the storage write path in a breaker: open after a failure-rate threshold over a rolling window (for example, 50% errors over the last 10 attempts), during which writes fail fast into the buffer instead of blocking on a slow database. Move to half-open after a backoff period to test recovery with a small amount of traffic before fully closing again.
Retry strategy
Exponential backoff with jitter for transient errors, with a hard cap so retries do not themselves become a load source:
giving delays of 200, 400, 800, 1600, 3200, 6400 ms for attempts 1 through 6, reaching the 30-second cap by around attempt 9.
Surfacing pipeline health
Expose queue depth, drop rate, breaker state, and write-success rate as first-class metrics with their own dashboard and alerting, plus a documented behavior contract (how long buffered data survives, what gets dropped first under sustained pressure) so dependent teams know what to expect during an incident instead of just seeing gaps.
Worked Example
Assume normal load into the time-series database is 50,000 points/sec, matched by normal write capacity, so no backlog accumulates in steady state. During an incident, traffic spikes 10x to 500,000 points/sec while write capacity simultaneously drops to 20% of normal (10,000 points/sec) because the same incident is degrading the database itself, a stated worst-case scenario for sizing purposes.
Backlog growth rate:
500,000−10,000=490,000 points/sSizing the buffer to survive a 5-minute (300 s) incident before either recovery or an operator decision:
490,000×300=147,000,000 pointsAt 150 bytes/point (consistent with the same encoded-point-size assumption used for ingest sizing elsewhere):
147,000,000×150 B=22.05×109 B≈22.05 GBProvisioning roughly 22 GB of buffer capacity is what actually backs a "we survive a 5-minute downstream outage at 10x traffic" claim. Set a shedding high-watermark at 80% of that: 0.8×22.05≈17.6 GB, past which the pipeline starts dropping lowest-priority series (non-alerting, debug-tier metrics) to preserve budget for anything feeding an active SLO or alert.
Trade-offs and Pitfalls
Buffering trades data loss for latency and cost: a bigger buffer survives a longer outage without dropping anything, but costs more to provision and, if it fills anyway, the operator now has a large backlog to drain (see S3's drain-time math for why backlog drain time can badly exceed outage length) rather than a clean, immediate failure.
A circuit breaker that opens too aggressively (a low failure threshold, a short window) can trip on ordinary transient blips and start buffering unnecessarily, adding latency for no real benefit; one that opens too conservatively keeps hammering an already-struggling database and can make the underlying incident worse. Tune the threshold against the sink's actual recovery behavior, not a default value copied from an unrelated system.
Shedding low-priority series under pressure only works if "priority" was decided in advance, not improvised during the incident. If every team believes their metrics are the important ones, the shedding policy needs an actual, pre-agreed tier assignment, or it becomes a political argument during the worst possible moment to have one.
Design the aggregation and partitioning strategy for a horizontally scalable time-series database that needs efficient single-metric queries at long retention. How would you choose sharding keys (metric name versus specific label sets), handle replication, and separate the read path from the write path, while avoiding hot shards?
Sample Answer
Direct answer
Shard primarily by metric name so a single-metric query only ever touches one shard, but split the shard's key space with a secondary hash over the label set so one enormously popular metric doesn't overload the shard that owns it. Separate the write path (append-heavy, needs low-latency durable ingestion) from the read path (needs efficient range scans and can tolerate a small replication lag) by giving reads their own replica tier, and detect hot shards from live load metrics rather than assuming the initial hash assignment will stay balanced as traffic shifts.
Structured elaboration
Primary sharding key: metric name. shard=hash(metric_name)modN. This keeps all time-series for one metric co-located, so "give me this metric's values" is a single-shard query instead of a scatter-gather across the whole cluster, which is the dominant query pattern this design optimizes for.
Secondary partitioning: label set, within a metric. Inside the shard(s) a metric owns, further split by hash(sorted_label_pairs) into K sub-partitions. This only matters for metrics whose write volume alone would overload a single shard; low-traffic metrics don't need it.
Replication: replication factor 3 across availability zones, quorum writes for durability (a write acknowledges once a majority of replicas confirm), asynchronous replication to read replicas that can serve slightly stale reads for range queries that don't require strict recency.
Read/write path separation: writes go to the shard leader's write-ahead log and in-memory segment; reads for recent data can be served from the leader or a follower depending on the caller's consistency requirement, and reads for long-range historical queries hit compacted, downsampled cold storage rather than the hot write path at all.
Rebalancing: use virtual nodes (many small virtual shards mapped onto fewer physical nodes) so that splitting a hot shard means moving a subset of its virtual nodes to other physical nodes, not re-hashing the whole ring.
flowchart LR
W[Writes] --> L[Shard leader: WAL + in-memory segment]
L -- quorum ack --> F1[Follower replica]
L -- quorum ack --> F2[Follower replica]
R[Range / historical reads] --> C[(Compacted cold storage)]
R2[Recent reads] --> F1
Q[Single-metric query] --> RT["Ring lookup: hash(metric_name)"]
RT --> L
Worked example
Assume 10,000,000 writes/sec cluster-wide across N=64 shards. If load were perfectly uniform, average shard load would be:
avgShardLoad=6410,000,000=156,250 writes/secSharding by metric name alone. Suppose one metric name (a widely-instrumented per-request counter) carries 8% of all cluster writes:
hotMetricWrites=0.08×10,000,000=800,000 writes/secBecause that entire metric maps to one shard under pure metric-name hashing, that shard's load is roughly the hot metric's full volume plus its normal share of other metrics:
hotShardLoad≈800,000+156,250=956,250 writes/sec avgShardLoadhotShardLoad=156,250956,250=6.12× the average shardThat shard is running over 6x hotter than every other shard in the cluster, purely from metric-name hashing.
Adding secondary label-hash partitioning (K=16 sub-partitions for this one metric):
perSubPartitionWrites=16800,000=50,000 writes/sec shardLoadAfterSplit≈50,000+156,250=206,250 writes/sec avgShardLoadshardLoadAfterSplit=156,250206,250=1.32×Splitting the hot metric's writes across 16 label-hashed sub-partitions brings its owning shards from 6.12x average load down to 1.32x, close enough to average that normal rebalancing headroom absorbs it. The insight worth stating explicitly: the number of sub-partitions K should scale with a metric's observed write share, not be a fixed constant applied uniformly, since a metric carrying 0.1% of traffic never needed splitting in the first place.
Trade-offs & pitfalls
| Sharding strategy | Single-metric query | Hot-metric protection | Rebalance cost |
|---|---|---|---|
| Metric name only | 1 shard, fast | None: any popular metric is a hot shard | Low, rare |
| Label set only | Scatter-gather across shards | Even by construction | Low |
| Metric name + secondary label split (this design) | 1 shard for low-traffic metrics, bounded fan-out for hot ones | Good, tunable per metric | Moderate: needs live monitoring to decide K |
Common wrong turns: applying secondary label-set splitting to every metric uniformly, which turns even low-traffic single-metric queries into unnecessary scatter-gathers across K sub-partitions for no load-balancing benefit; treating quorum-write durability and hot-shard mitigation as unrelated concerns when a hot shard under quorum writes also means its followers are under proportionally more replication load, compounding the imbalance; and rebalancing by re-hashing the whole ring on every detected hot shard instead of moving individual virtual nodes, which causes a much larger and slower data migration than the imbalance warranted.
You need to deploy an OpenTelemetry Collector fleet that can autoscale with load and keep accepting data even if the downstream backend has an outage. How would you design the deployment (agent versus gateway, horizontal autoscaling, a durable buffer sitting in front of the exporters) and structure the processor chain, for example batching, sampling, and enrichment?
Sample Answer
Direct Answer
Split the fleet into a fixed agent tier (one per host, DaemonSet-deployed, doing only light initial processing) and a horizontally autoscaled gateway tier that does the heavier work (sampling, enrichment, batching) and fronts the actual export. Put a durable buffer between the gateway and the exporter so a backend outage fills the buffer instead of blocking or dropping at the gateway, and scale the gateway tier on real backpressure signals (queue depth, CPU) rather than a fixed replica count.
Structured Elaboration
Fleet topology
flowchart LR
AGENT["Per-Host Agents"] --> GW["Gateway Tier: HPA-scaled"]
GW --> PROC["Processor Chain: memory_limiter, enrich, sample, batch"]
PROC --> BUF[("Durable Buffer")]
BUF --> EXP["Exporter"]
EXP --> BACKEND[("Backend")]
BUF -.->|"queue depth metric"| HPA{"Autoscaler"}
HPA -.->|"scale replicas"| GW
Agent versus gateway split
Agents run one per host, receiving OTLP from local processes and forwarding onward with minimal processing, so their resource footprint per host stays flat and predictable regardless of fleet-wide load. The gateway tier does the load-dependent work (see below), which is exactly what needs to scale with traffic, so isolating it from the fixed agent tier is what makes autoscaling meaningful.
Horizontal autoscaling
Scale the gateway tier on a metric that reflects real backpressure, most reliably queue depth in front of the exporter (or CPU as a proxy if queue depth isn't exposed as a scalable metric), using a standard proportional scaling rule: desired replica count scales with how far the current metric is above target.
Durable buffer in front of the exporters
A persistent queue (backed by local disk or an external durable log like Kafka) sits between the processor chain and the exporter. During a backend outage, the gateway keeps accepting and processing data, writing into the buffer instead of failing the export, and drains the buffer once the backend recovers.
Processor chain structure
Order matters: memory_limiter first (shed load before spending CPU on anything else), then resource and attribute enrichment, then sampling (tail-based sampling needs to see whole traces, so it runs before batching, not after), then batch last, right before the exporter, so batches only ever contain data that already survived the sampling decision.
Worked Example
Autoscaling thresholds. Assume each gateway pod sustains 20,000 spans/sec (a stated sizing assumption reflecting the processor chain's per-pod cost), and incoming load ranges from a floor of 50,000 spans/sec to a peak of 400,000 spans/sec through the day.
minReplicas=⌈20,00050,000⌉=3,maxReplicas=⌈20,000400,000⌉=20Applying the standard horizontal-scaling proportional rule, desiredReplicas = ceil(currentReplicas x currentMetricValue / desiredMetricValue), for a concrete example at 5 current replicas running at 90% CPU against a 70% target:
desiredReplicas=⌈5×7090⌉=⌈6.43⌉=7Durable buffer sizing. Assume the gateway continues accepting spans at 100,000/sec during a backend outage (average 800 bytes/span post-enrichment, a stated design input), and the outage is expected to last up to 15 minutes (900 s) before either recovery or a paging escalation forces a decision:
100,000×900×800 B=7.2×1010 B=72 GBA roughly 72 GB buffer requirement is well within the local disk capacity of a typical commodity broker or persistent-volume-backed queue, so this is a provisioning number to size the buffer's storage class against, not a hard architectural constraint.
Trade-offs and Pitfalls
Scaling on CPU alone is a lagging proxy: CPU can look fine even while the queue in front of the exporter is growing, if the bottleneck is actually export throughput to a slow backend rather than processing throughput. Scaling on queue depth directly is a better signal when it's available, since it reflects the thing you actually care about (are we falling behind) rather than a correlated but imperfect stand-in.
A durable buffer removes the immediate pressure to fix an outage fast, which is good for availability but can mask a slow backend degradation if nobody is watching buffer fill rate: a buffer that's growing steadily but not yet full looks the same as a healthy system on a dashboard that only shows "buffer not full yet," unless fill rate itself is alerted on.
Running sampling before batching (rather than after) is correct for tail-based decisions that need whole traces, but it means the sampling stage has to hold and correlate spans across a trace, which is real memory pressure on the gateway tier specifically, on top of whatever the autoscaling math above accounts for from raw throughput alone. Size gateway pod memory against expected in-flight trace count, not just span throughput.
Unlock Full Question Bank
Get access to all 40 Observability and Monitoring Architecture interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.