Problem Statement: A Planet-Scale Log Ingestion and Alerting Platform
Frames the system as a high-throughput observability backbone, not merely a log viewer.
Problem statement
Design a distributed logging and alerting system that collects structured and unstructured log events from thousands of microservices spread across multiple data centers and cloud regions, ships them through a durable ingestion pipeline, indexes them for sub-second full-text and structured search, evaluates threshold-based and pattern-based alert rules in near real time, correlates related events into incidents, and routes actionable notifications to on-call engineers through configurable escalation policies.
This is not a simple tail -f aggregated into a dashboard. At the scale of a large internet company, the system must absorb millions of log events per second, tolerate the loss of any single ingestion node without data loss, provide query latency under two seconds for recent data and under ten seconds for historical data spanning thirty days, and evaluate tens of thousands of alert rules continuously with a false-negative rate below 0.1 percent.
Why this problem is distinctive
A web application can retry a failed database write. A logging pipeline cannot retry a log line that was never captured: the process that emitted it may have already crashed, the container may have been recycled, and the ephemeral context is gone forever. The design therefore separates ingestion durability from query performance and alert accuracy. Ingestion durability is a correctness invariant: no log event accepted by an agent may be silently dropped before it reaches durable storage. Query performance is an optimization target: indexes, caches, and tiered storage trade latency against cost. Alert accuracy is a signal-quality target: precision and recall must both be high, because a flood of false positives trains engineers to ignore the system, and a missed true positive means a production outage goes undetected.
The four architectural planes
- Collection plane: lightweight agents running on every host and sidecar, tailing log files, capturing stdout/stderr, and shipping structured events with backpressure.
- Transport plane: a durable, partitioned, replayable message backbone that decouples producers from consumers and absorbs traffic spikes.
- Processing and storage plane: stream processors that parse, enrich, route, and index events into search-optimized stores with tiered retention.
- Alerting and correlation plane: a rule evaluation engine, an event correlator, an incident manager, and a notification router with escalation policies.
A strong interview answer keeps these planes separate. It allows the storage plane to degrade query latency without losing ingestion durability, and it permits the alerting plane to sample or approximate when full-volume evaluation exceeds capacity, without silently dropping events from the durable log.
Public operating baseline versus design assumptions
Public evidence establishes that the category is operationally real and massive. LinkedIn engineering reports that Apache Kafka, originally built for LinkedIn's activity streams and log aggregation, handles over 7 trillion messages per day across multiple clusters. Uber's engineering blog describes a logging pipeline that processes over 1 trillion events per day from thousands of microservices. Datadog's S-1 filing and investor materials report ingesting over 30 trillion log events per month across their customer base. Grafana Labs reports that Loki, their cloud-native log system, is used in production by thousands of organizations and processes petabytes of log data.
For capacity planning, this answer explicitly assumes a mature enterprise platform with 15,000 services, 200,000 containers, 5 million log events per second at average load, and a five-times event peak. Unless a number is tied to a citation, it is a stated design assumption, target, budget, or illustrative threshold—not a claim about any company's private architecture.
Scope boundaries
We design the end-to-end platform: agent collection, transport, parsing, enrichment, indexing, search, alert rule management, evaluation, correlation, notification routing, retention, and multi-region replication. We do not design the application logging SDKs themselves (we assume structured logging libraries exist), the underlying container orchestrator, or the identity provider. We still design the interfaces, contracts, and failure semantics around those dependencies.
Key Highlights
- •The cloud may optimize search and alerting, but ingestion durability is a local invariant: no accepted event is silently lost.
- •Model ingestion as a durable pipeline, search as a tiered optimization, and alerting as a signal-quality problem.
- •Public fleet figures from LinkedIn, Uber, and Datadog provide context; every uncited scale or SLO in this answer is an explicit design assumption.
- •The architecture has four planes: collection, transport, processing/storage, and alerting/correlation.
- •A missed alert is a safety failure; a false alarm flood is a trust failure. Both must be designed against.
Section Rescue Kit
Buzzwords to use:
Safe statements:
- "I will separate ingestion durability from query performance and alert accuracy because they have different consistency and failure requirements."
- "Before selecting technologies, let me define which guarantees are hard invariants and which are optimization targets."