Interview Setup
Interview Prompt
Design a log aggregation and search system like Splunk or the ELK stack. Ingest from thousands of services, search with under 5 seconds latency, stream live tails, and alert on error patterns.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| What is the ingest rate and retention window? | 1M lines/sec at 30 day retention produces about 130 TB of compressed logical data before index and replica overhead, making storage tiering a major design decision. |
| Full text search on message bodies or label based filtering only? | Elasticsearch indexes message content for full text search, which can materially increase stored index footprint, while Loki indexes labels and uses them to narrow the set of compressed chunks that must be scanned. Treat the ten-fold storage reduction as an illustrative workload assumption rather than a universal product ratio. |
| How fresh must search results be after ingestion? | Near real time availability within seconds requires bulk indexing with short refresh intervals, whereas batch indexing over several minutes significantly reduces operating costs. |
| Do we need live tail, or is search only sufficient? | Live tail requires WebSocket fan-out and shared filtered Kafka consumers, which creates a separate scaling problem from historical search. |
Scope
In scope
- Agent to Kafka to indexer pipeline
- Structured logging
- Index rotation and retention
- Full text search at scale
- Hot, warm, and cold tiering
- Log based alerting
Out of scope (state explicitly)
- Application instrumentation SDK design
- Full distributed tracing system (explored in Distributed Tracing)
- On call paging and escalation policy (explored in On Call Escalation System)
Functional Requirements
Log collection requires centralized ingestion across heterogeneous fleets, high speed full text indexing, retention tiers, and structured field extraction, while balancing real time tailing against historical investigation queries.
This architecture covers the log ingestion pipeline using Elasticsearch as a sink, incorporating log agents, Kafka message buffering, and lifecycle tiers. For operating an Elasticsearch cluster as a dedicated product, including cluster roles, shard topologies, and scatter-gather execution, see the Elasticsearch Search Cluster guide.
- Log Ingestion: Collect and ingest logs continuously from thousands of microservices, physical hosts, container clusters, and cloud platform components.
- Format Support: Parse and normalize structured and unstructured log payloads, including JSON, syslog, raw text, and multiline exception stack traces.
- Full Text Search: Execute substring and full text searches across all ingested log events with under 5 seconds query latency across recent data.
- Multidimensional Filtering: Filter event streams dynamically by service name, severity level, time window, host identifier, and custom structured attributes.
- Live Tail Streaming: Stream matching log events to connected client consoles in real time over persistent connections.
- Aggregation and Dashboards: Aggregate volume metrics, error rates, and operational trends for visualization across real time dashboards.
- Automated Alerting: Trigger incident alerts when error log frequencies breach configured thresholds or when critical error patterns emerge.
- Configurable Retention: Enforce data retention policies per log source across a lifecycle spanning 7 days to 1 year.
Non-Functional Requirements
The platform must ingest tens of terabytes daily without dropping acknowledged logs during sudden traffic bursts, while returning bounded recent-data search results within the under 5 second target.
- High Throughput: Ingest over 1,000,000 log lines per second sustained (translating to 500 MB/s sustained ingestion bandwidth).
- Low Latency Search: Return query results in under 5 seconds across logs generated in the preceding 24 hours.
- Scalability: Scale horizontally to store petabytes of historical logs originating from thousands of distinct services.
- Durability: Guarantee zero log loss after the local collection agent durably writes an event to its persistent buffer and acknowledges it. If the durable buffer reaches capacity, the agent must stop acknowledging additional events or apply an explicitly documented loss policy to unacknowledged data.
- Cost Efficiency: Apply aggressive compression and automated hot, warm, and cold storage tiering to control infrastructure spend.
- High Availability: Provide 99.99% uptime for both log ingestion and active query services.
Capacity Estimations
Daily ingest volume and retention duration directly determine the sizing of hot, warm, and cold storage tiers because physical input/output throughput and storage footprints dominate costs at scale. Sustaining 1,000,000 log lines per second requires aggressive compression and automated lifecycle tiering to maintain predictable operating expenses. The one-year retention calculation is based on 4.3 TB/day of compressed logical data, while the cost example applies the stated storage prices to that compressed volume and excludes index, replica, compute, request, and egress costs.
| Metric | Calculation | Value |
|---|---|---|
| Log lines / sec | Scenario peak ingestion rate | 1M |
| Avg log line size | Typical production workload assumption | 500 bytes |
| Ingestion throughput | 1M lines/sec x 500 bytes/line | 500 MB/s |
| Storage / day (raw) | 500 MB/s x 86,400 seconds | 43 TB |
| With compression (~10x) | 43 TB/day raw ÷ 10x compression ratio | 4.3 TB/day |
| 30 day retention | 4.3 TB/day x 30 days | ~130 TB |
| 1 year compressed retention | 4.3 TB/day x 365 days | ~1.57 PB before index and replica overhead |
| 30 day active hot and warm footprint | 4.3 TB/day x 30 days | ~130 TB before index and replica overhead |
Architecture Diagram
Logs travel from host level collection agents into a distributed message buffer, proceed through stream processors into a search indexing cluster, and are accessed via dedicated query and live tail APIs. Because storage footprint and I/O throughput dominate operating expenses at scale, the architecture pairs high performance SSD hot tiers with aggressive compression and progressive migration to object storage.
Component Deep Dives
Event Bus Design (Kafka)
The ingestion path is fully decoupled from search. Local agents handle host level concerns such as buffering, multiline stack trace assembly, and PII redaction, while Kafka and downstream indexer workers absorb volume bursts without creating backpressure on the applications that emit logs.
Kafka acts as the primary absorption shock absorber between thousands of log generating producers and the search indexing cluster.
Topic: logs-raw
Partitions: 1,024
Partition key: hash(service_name, source_id) (spreads a busy service across sources while preserving per-source ordering)
Retention: 48 hours (buffer for reprocessing and live tail)
Replication factor: 3, min.insync.replicas: 2
Producer durability: acks=all, enable.idempotence=true, retries enabled
Broker safety: unclean.leader.election.enable=false (prefer availability loss over acknowledged data loss)
Producer: Log Agent (Vector or Filebeat) after parse, enrich, PII redaction, and compress stages
Event: { timestamp, event_id, service_name, source_id, host, level, message, trace_id, labels }
event_id: stable unique identifier assigned once at collection and retained across retries and replays
Consumer groups:
1. log-processor: deep parse, PII validation, and bulk index into Elasticsearch
5,000 docs per 5 seconds per worker equivalent
Elasticsearch document ID = event_id for idempotent indexing
Retry failed bulk items before offset acknowledgement
Permanently failed records go to logs-raw-dlq before the source offset is acknowledged
2. alert-evaluator: pattern match for critical alerts (OOM, disk full, error rate) and forward to PagerDuty or Slack
3. live-tail: shared consumer pool reads a bounded set of filtered topics and routes matching events to WebSocket sessions through a subscription registry, with event_id deduplication at the routing layer
Capacity note: 5,000 docs per 5 seconds is 1,000 docs/sec per worker equivalent. At 1M log lines/sec, the system needs approximately 1,000 effective worker equivalents or higher per-worker throughput, so Kafka partitions and indexer capacity must scale together. The 1,024 partition configuration leaves enough source parallelism for this illustrative target, with additional headroom provided by faster workers or future partition expansion.
Ingestion path: Agent local buffer to Kafka to stream processor to Elasticsearch hot tier (0 to 48 hours)
Lifecycle tiering: Hot SSD cluster to warm HDD cluster to cold S3 Parquet to Glacier archive
Dead-letter queue: logs-raw-dlq for records that fail processing after bounded retries. Alert on DLQ growth separately from consumer lag, and alert on log-processor lag exceeding 60 secondsAPI Design
Log Search and Live Tail Interfaces
The log aggregation platform exposes search query endpoints with precise time range filters, alongside persistent WebSocket connections for live log streaming. Each immutable log event carries a stable event ID so downstream indexing and live tail routing can safely deduplicate retries.
type LogTimestamp = string; // ISO-8601 UTC timestamp
type LogCursor = string; // Opaque cursor containing search_after and optional PIT state
type LogQuery = string;
type LogLevel = "DEBUG" | "INFO" | "WARN" | "ERROR" | "CRITICAL";
type LogFieldValue = string | number | boolean;
type LogSort = "timestamp:desc,event_id:desc" | "timestamp:asc,event_id:asc";
type LogFilter = Record<string, LogFieldValue>;
export interface LogSearchRequest {
query: LogQuery;
fromTimestamp: LogTimestamp;
toTimestamp: LogTimestamp;
limit?: number;
sort?: LogSort;
filterFields?: LogFilter;
cursor?: LogCursor;
}
export interface LogEntry {
timestamp: LogTimestamp;
service: string;
host: string;
level: LogLevel;
message: string;
event_id: string;
trace_id?: string;
span_id?: string;
labels?: Record<string, string>;
}
export interface LogSearchResponse {
totalHits: number;
tookMs: number;
logs: LogEntry[];
nextCursor?: LogCursor;
}
export interface LiveTailSubscription {
type: "subscribe" | "unsubscribe";
filter: string;
maxLinesPerSecond?: number;
}
export interface LiveTailMessage {
type: "log" | "heartbeat" | "error";
data?: LogEntry;
message?: string;
}Search Logs Endpoint
Queries indexed logs using full text search syntax, supporting Boolean operators and timestamp filters. Wildcard searches use a dedicated wildcard field and should be bounded by time range and result limits because broad leading wildcard searches can be expensive.
POST /api/v1/logs/search
Content-Type: application/json
{
"query": "level:ERROR AND service:payment-service AND message_wildcard:*timeout*",
"from": "2026-03-14T04:00:00Z",
"to": "2026-03-14T10:00:00Z",
"limit": 100,
"sort": "timestamp:desc,event_id:desc"
}Deep Pagination
The nextCursor should encapsulate Elasticsearch search_after sort values rather than relying on deep from and size offsets. The sort includes event_id as a deterministic tie breaker. For pagination that must remain consistent while the index refreshes, the service can pair search_after with a point in time snapshot and return the snapshot identifier inside the cursor. The point in time context should have a bounded lifetime, and an expired cursor should return a clear pagination error so the client can restart the query.
Live Tail (WebSocket)
Establishes a persistent bidirectional connection to stream matching log events to the client in real time.
Client subscription message:
{
"type": "subscribe",
"filter": "service=payment-service AND level=ERROR"
}Server event message:
{
"type": "log",
"data": {
"timestamp": "2026-03-14T10:05:00.123Z",
"service": "payment-service",
"level": "ERROR",
"message": "Gateway timeout: Stripe API failed to respond",
"event_id": "evt_99f2b8a_170",
"trace_id": "trace_99f2b8a",
"host": "payment-pod-84b"
}
}Common Error Responses
400 Bad Request: invalid input, missing required fields, or malformed JSON payload 401 Unauthorized: missing or invalid authentication token or API key 403 Forbidden: authenticated caller lacks required permissions for this resource 404 Not Found: requested resource ID does not exist 409 Conflict: duplicate write or version conflict, retry with a unique idempotency key 422 Unprocessable Entity: syntactically valid request failed semantic business validation 429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header 500 Internal Error: unexpected server failure, retry safely with an idempotency key 503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff 504 Gateway Timeout: search index shard responded slowly, narrow query parameters or retry
Data Model
Elasticsearch Index Mapping
Structured logs map to specific field data types to optimize disk footprint and query performance. Exact-match and aggregation fields such as service, host, level, and trace_id can use keyword mappings, while arbitrary custom attributes use flattened to avoid uncontrolled dynamic field growth. The message_wildcard field is dedicated to grep-like wildcard queries so the main full text field remains optimized for analyzed search.
{
"mappings": {
"properties": {
"timestamp": {"type": "date"},
"level": {"type": "keyword"},
"service": {"type": "keyword"},
"host": {"type": "keyword"},
"env": {"type": "keyword"},
"message": {"type": "text", "analyzer": "standard", "copy_to": "message_wildcard"},
"message_wildcard": {"type": "wildcard"},
"event_id": {"type": "keyword"},
"trace_id": {"type": "keyword"},
"custom_fields": {"type": "flattened"}
}
},
"settings": {
"index.number_of_shards": 5,
"index.number_of_replicas": 1,
"index.lifecycle.name": "logs-policy"
}
}Index Lifecycle Management (ILM) Policy
Automates the progression of indices down the physical storage hierarchy according to log age and access patterns:
# Illustrative lifecycle model, not a direct Elasticsearch ILM API payload
# Elasticsearch ILM uses hot, warm, cold, frozen, and delete phases.
# Long term Glacier archival is handled through the snapshot/archive workflow.
hot_phase:
retention: "0 to 2 days"
storage_media: "NVMe SSD"
sharding: "5 primary shards, 1 replica"
actions: "Force merge after rollover when the index is no longer actively written"
warm_phase:
retention: "2 to 30 days"
storage_media: "Attached HDD"
sharding: "Shrink to 1 primary shard, read only"
actions: "Best-effort block compression"
cold_phase:
retention: "30 to 90 days"
storage_media: "Object Storage (S3 Standard-IA)"
sharding: "Searchable snapshot, 0 replicas"
actions: "Searchable snapshot mount"
archive_phase:
retention: "90 days to 1 year"
storage_media: "Glacier archive"
actions: "Archive after searchable retention ends"
delete_phase:
retention: "Older than 1 year"
storage_media: "Permanent deletion after policy retention"
actions: "Purge index metadata and delete underlying archived data"Fault Tolerance
Log Volume Spikes and Cascading Retries
The ingestion pipeline must handle physical bottlenecks, downstream system failures, and unexpected log volume surges.
| Concern | Solution |
|---|---|
| Agent failure | Agent uses a persistent local disk buffer that survives process restart and resumes shipping after recovery. Host loss requires replicated remote buffering for the same durability guarantee |
| Kafka down | Each agent buffers locally up to 1 GB, retries with backoff, and stops accepting events into the durable buffer once that per agent limit is reached, or follows the explicitly configured overflow policy. The 1 GB limit is per agent, not an aggregate platform buffer |
| Elasticsearch overload | Kafka absorbs bursts while indexer workers scale horizontally and query load is rate limited |
| Index corruption | Elasticsearch replicas help with node or shard failures, while daily S3 snapshots provide a recovery point for broader corruption. Tested restore procedures are required, and recent data after the latest snapshot may still need replay from Kafka |
| Query overload | Query timeout set at 30 seconds alongside per user rate limiting and bounded search ranges |
| Disk full | ILM enforces retention, alerts when disk usage exceeds 80%, and emergency deletion is limited to data already beyond policy |
When an upstream dependency fails (such as a database connection timeout), hundreds of microservice pods may log errors, retry immediately, and generate duplicate log entries. In a cluster of 500 pods, this failure cascade can suddenly flood the logging platform with 50,000,000 log lines per second (50x the 1,000,000 line per sec baseline), threatening to exhaust Kafka disk space, push consumer lag to multiple hours, and run search cluster nodes out of storage.
- Agent Level Rate Limiting: The local log daemon caps outbound shipping at 1,000 log lines per second per container. Excess logs are sampled or dropped with an aggregated periodic summary line recording dropped counts.
- Kafka Per Service Quotas: Broker quotas throttle misbehaving client producers sending disproportionate traffic, propagating backpressure upstream to the application logging framework.
- Circuit Breaking in Log Appenders: If the logging transport blocks or drops frames, asynchronous log appenders drop entries gracefully or convert them to lightweight telemetry counters.
- Tiered Ingestion via Severity Shedding: During ingestion emergencies, indexer workers discard
DEBUGandINFOlogs to guarantee ingestion capacity forERRORandCRITICALevents.
Multiline Stack Trace Assembly
Multiline runtime exceptions, such as Java or Python stack traces, emit separate lines to standard output. If ingested naively line by line, a single exception splits into multiple distinct log entries, destroying search context and corrupting trace correlations.
- Agent Level Pattern Buffering: The collection agent buffers multiline blocks using regex continuation patterns (such as lines starting with whitespace,
Caused by, orat). - Flushing Timeouts: If no new continuation lines arrive within 500 milliseconds, the agent flushes the buffered block immediately to keep delivery latencies minimal.
- Structured JSON Logging: Forcing applications to serialize exception stack traces directly into a structured JSON string on a single line eliminates multiline regex assembly entirely.
Elasticsearch Shard Explosion and Heap Pressure
An infrastructure ingesting 4.3 TB of compressed logs daily across 50 distinct services can create severe shard growth if every service creates daily indices. With 5 primary shards and 1 replica, retaining 30 days of daily indices creates roughly 50 x 30 x 5 x 2 = 15,000 shard copies. High active shard counts increase cluster-state size, coordination work, and master JVM heap pressure.
- Unified Data Streams: Group smaller services into shared time based data streams rather than dedicated per service indices, filtering by the indexed
servicekeyword at query time. - Rollover Force Merging: Shrink 5 primary shards down to a single read only shard during warm rollover, reducing Lucene segment overhead and open file descriptors.
- Replica Pruning: Drop read replicas on indices older than 7 days once write traffic ceases only when the retained snapshot and recovery policy provides the required durability and recovery point objective.
Live Tail Streaming at High Scale
Users observing live tails expect instantaneous log output in their console. However, polling Elasticsearch every second causes severe cluster query load, while streaming directly from the raw 1,000,000 lines/sec Kafka topic for every active user session exhausts consumer bandwidth.
- Pre Filtered Kafka Topics: Dedicated stream processors route high frequency matching patterns, such as
ERRORevents per service, to a bounded set of secondary filtered Kafka topics. Do not create one Kafka topic per live tail session. - Distributed WebSocket Session Filtering: Cap concurrent live tail sessions at 50 active users per cluster, falling back to 5 second polling against indexed storage beyond this illustrative threshold to safeguard ingestion stability.
Additional Considerations
PII Redaction: Preventing Sensitive Data in Logs
Developers frequently log sensitive customer details by accident, including email addresses, authentication secrets, Social Security numbers, and credit card credentials. Storing sensitive data in plaintext can violate security and privacy requirements such as GDPR and CCPA. Retrospective cleanup across petabytes of indexed data can be expensive and slow, so prevention at the source is mandatory.
Defense in Depth Solution:
- Layer 1 (Code Linting): Reduce accidental PII logging through CI rules that flag typical sensitive parameter names and known logging anti-patterns.
- Layer 2 (Agent Redaction): Local agent executes high speed regex sanitization matches:
- Email regex:
[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,} - SSN regex:
\d{3}-\d{2}-\d{4} - Credit Card regex:
\d{4}[- ]?\d{4}[- ]?\d{4}[- ]?\d{4}
- Email regex:
- Layer 3 (Tokenization Vaults): Replace identified PII parameters with reversible transaction tokens (such as
tok_39d8c4). The tokenization service controls key access through managed key storage and role based authorization, so individual operators do not hold the underlying cryptographic keys.
Access Control and Tenant Isolation
Logs often contain sensitive production data, so search and live tail access must be authenticated and authorized independently of the ingestion path. The service injects tenant, environment, and service scope from the authenticated identity rather than trusting equivalent values supplied only in the query, and it audits sensitive log access. Encryption in transit and at rest is mandatory, while agent side redaction remains the first line of defense for secrets and regulated fields.
Structured vs Unstructured Logging: Why JSON Wins
Unstructured logs require complex and fragile Grok or regular expression processing pipelines that break easily when log formats evolve:
Unstructured raw log line:
2026-03-14 10:00:00 ERROR [payment-service] User 12345 payment failed: timeout after 5000ms
Structured JSON log line:
{"timestamp":"2026-03-14T10:00:00Z","level":"ERROR","service":"payment-service","user_id":"12345","message":"payment failed","event_id":"evt_abc123","error":"timeout","duration_ms":5000,"trace_id":"abc123"}Why JSON wins:
- No application specific regex parsing: Structured fields avoid runtime Grok or regex extraction in the indexer for fields that are already typed, while JSON parsing itself remains necessary.
- Type preservation: Numeric values (such as
duration_ms) remain queryable numbers rather than strings. - Compression friendly: JSON structures repeating identical field keys are highly deduplicated by block compression algorithms like LZ4 and ZSTD.
Log, Trace, and Metric Correlation: The Observability Triangle
Tracing a user request failure across multiple microservices is challenging without system-wide correlation.
Unified Troubleshooting Workflow:
- The application SDK injects
trace_idandspan_idfrom OpenTelemetry tracing contexts into every structured log event payload. - When an error log alert triggers, engineers click the attached
trace_idto open a distributed waterfall trace view. - After pinpointing the specific database call causing latency, responders navigate directly to the database metrics dashboard for root cause analysis.
{
"level": "ERROR",
"message": "payment timeout",
"trace_id": "abc123",
"span_id": "def456"
}Cost Optimization: Calculations & Levers
At 4.3 TB/day compressed logical volume, keeping one year of data on hot SSD at the stated $0.10/GB-month price costs approximately $1.88M / year before index and replica overhead. Applying an aggressive tiering lifecycle architecture reduces the storage-only estimate to approximately $451K / year under the stated illustrative prices:
Annual Storage Cost Model: Assumptions: - 4.3 TB/day compressed logical log volume - Storage prices are stated per GB-month - Excludes Elasticsearch index overhead, replicas, compute, requests, and egress - Hot tier (2 days): 4.3 x 2 x 1000 x $0.10 x 12 = $10,320/year - Warm tier (28 days): 4.3 x 28 x 1000 x $0.03 x 12 = $43,344/year - Cold tier (335 days): 4.3 x 335 x 1000 x $0.023 x 12 = $397,578/year ------------------------------------------------------------------------- Tiered total: ~$451,242/year Pure hot storage: 4.3 x 365 x 1000 x $0.10 x 12 = ~$1,883,400/year Estimated reduction: ~76% through architectural tiering.
Additional Cost Levers
- Intelligent Sampling: Retain 100% of
ERRORandWARNlogs, but sample only 10% of high volumeINFOandDEBUGlines. This cuts total volume by approximately 70% only when INFO and DEBUG together account for about 77.8% of baseline traffic, so the actual reduction must be measured from production mix. - Heartbeat Filtering: Filter out continuous, redundant health check probe logs at the host agent level before they traverse the network.
OpenTelemetry Logs & Trace Correlation
Shipping logs with trace_id and span_id fields via the OpenTelemetry log bridge allows engineers to jump seamlessly from an individual log event to the full distributed transaction trace explored in Distributed Tracing. Sampling verbose DEBUG lines aggressively while retaining all ERROR events with complete context ensures that correlating logs to metric exemplars closes outage investigations significantly faster than search queries alone.
Related Problems and Core Concepts
Log aggregation and real time observability intersect directly with several platform design problems:
For deep-dive query execution, shard allocation, and cluster operations in search engines, explore Elasticsearch Search Cluster. Distributed trace propagation and context bridge mechanisms are analyzed in Distributed Tracing, while automated incident routing is detailed in On Call Escalation System. For time series telemetry and high throughput streaming pipelines, see Distributed Metrics Aggregation and User Analytics Pipeline.
To master the architectural concepts underpinning this design, review Kafka Architecture and Guarantees, Indexing and Query Optimization, Observability and Distributed Tracing, Stream Processing Basics, and System Design Interview Patterns.
Interview Walkthrough
- 25-minute cut
Skip arch50 and arch75 depth unless interviewing for a staff-level role.
- Diagram the pipeline from host agent to Kafka buffer to indexer workers and Elasticsearch, explaining each component's backpressure role (5 min)
- Quantify ingestion scale at 1,000,000 lines/sec and 4.3 TB/day compressed, justifying the Kafka buffer and ILM storage tiers (6 min)
- Mandate structured JSON logging with low cardinality index fields configured at the edge (5 min)
- Detail hot, warm, and cold ILM: active logs on SSD, aged indices on HDD, and long-term retention on object storage (5 min)
- Staff level: multiline stack trace assembly and edge PII redaction in the collector pipeline (4 min)
- Trace the ingestion path: application runtime to local host agent with buffering and redaction, through Kafka and indexer workers, down to hot, warm, and cold storage tiers.
- Mandate structured JSON logging across services because unstructured Grok parsing breaks on format drift and introduces parsing latency.
- Redact PII at the host agent layer using regex detection and replacement before logs leave the node, because retrospective cleanup across petabyte archives can be expensive and slow.
- Implement progressive storage tiering: 2 days on hot SSD Elasticsearch, 28 days on warm HDD, and 335 days on cold object storage for an estimated 76% reduction in storage cost under the stated pricing assumptions.
- Inject trace_id from OpenTelemetry contexts into every log record to enable one click navigation from error logs to distributed trace waterfalls and metrics.
- Apply intelligent sampling by retaining 10% of INFO and DEBUG logs while preserving 100% of ERROR and WARN events, reducing ingest volume by approximately 70% when INFO and DEBUG comprise about 77.8% of baseline traffic.
- Quantify scale: 500,000 to 1,000,000 logs/sec produces 500 MB/s to 1 GB/s ingestion bandwidth, requiring local agents to batch and compress before network transfer.
- Address the critical pitfall of retaining all logs in hot Elasticsearch indefinitely, which costs approximately $1.88M annually under the stated compressed volume and $0.10/GB-month storage assumption, before index and replica overhead.
Engineering Trade-offs
Storage Engine Architecture Comparison
Selecting the appropriate storage engine requires balancing full text query flexibility, ingestion throughput, and petabyte scale storage costs.
| Feature | Elasticsearch | Grafana Loki | ClickHouse |
|---|---|---|---|
| Index model | Full text inverted index on configured text fields | Index metadata labels only without log message indexing | Columnar store with secondary skip indices |
| Storage cost | High (illustrative index and segment footprint can reach 1.5x to 2x logical log volume) | Low (compressed chunks plus a small label index) | Medium (efficient columnar compression) |
| Query speed | Fast for arbitrary ad hoc full text queries | Slower for broad unstructured text queries because it narrows by labels, then scans matching chunks | Fast for structured aggregations and SQL filtering |
| Best for | Search heavy investigations and ad hoc troubleshooting | High volume ingest with label based operational filtering | Structured event streams and real time analytical SQL |
| Operational complexity | Complex (JVM tuning, master node state, shard management) | Moderate (distributed components backed by object storage) | Moderate to high (cluster management and Keeper or ZooKeeper coordination, depending on deployment) |
Indexing Everything vs Indexing Labels Only
The core architectural decision centers on whether to index every word across all message bodies into inverted indices, or to index only metadata labels while storing raw compressed message blocks.
Full Inverted Indexing (Elasticsearch): - Every word in every log line is indexed into an inverted index - Bounded wildcard search queries on the dedicated message_wildcard field can meet the under 5 second target for recent data - Ingestion CPU overhead is substantial during tokenization and segment creation - Illustrative storage assumption: index files consume 1.5x to 2x the logical log volume Label-Only Indexing (Grafana Loki): - Only metadata labels are indexed (such as service, level, host, and environment) - Log message bodies persist as raw compressed chunks without inverted index overhead - Illustrative workload assumption: achieves a 10x reduction in storage footprint and 10x higher ingestion throughput - Full text search narrows streams by labels, then scans matching compressed chunks, increasing query latency for broad time ranges
Decision Matrix. Select Grafana Loki when operational queries predominantly filter by known service tags, environments, or container metadata. Select Elasticsearch when engineering incident response requires arbitrary full text wildcard searches and complex distributed aggregations across the entire microservice fleet. High cardinality labels should be controlled explicitly in either design because unbounded label growth can increase memory, index, and query costs.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.