Interview Setup
Interview Prompt
Design a distributed tracing system (like Jaeger or Zipkin) that collects spans from microservices, assembles them into traces, supports head-based and tail-based sampling, and stores and queries traces efficiently for debugging production issues.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| What is the expected scale in terms of spans per second and retention window? | A 10M spans per second stress scenario with a 1% head-sampled mode produces 100K spans per second retained in that lower-cost mode, while the canonical full-fidelity path targets about 10% average retained volume. This demonstrates why the storage tier and indexing strategy must be decoupled from raw ingestion. |
| Should we optimize for head-based sampling, tail-based sampling, or an adaptive approach? | Head sampling introduces minimal overhead but misses rare anomalies, whereas tail sampling captures all errors at the cost of memory buffering. |
| What are the primary query patterns across trace identifiers, service names, and latency boundaries? | Direct trace ID lookups are bounded by the number of spans in the trace, while multi-dimensional range queries across services and latency boundaries require indexed columnar storage. |
| How should traces integrate with existing metrics and log aggregation platforms? | Metrics exemplars directly reference trace identifiers, while structured logs embed trace IDs to unify telemetry across the three observability pillars. |
Scope
In scope
- Span and trace data model
- Context propagation (W3C Trace Context)
- Head-based and tail-based sampling
- Trace assembly from distributed spans
- Columnar storage for span data
- Query by trace ID and service/latency
Out of scope (state explicitly)
- Low-level bytecode manipulation in APM language runtimes
- General metrics aggregation and log search systems (covered in Distributed Metrics Aggregation and Log Aggregation and Search). Trace-derived latency, error, and dependency aggregates needed by this tracing system remain in scope.
- Real-time anomaly detection and root-cause machine learning engines
Functional Requirements
Open by clarifying the expected observability depth with the interviewer. Essential capabilities include span instrumentation, pipeline collection, correlation into end-to-end traces, and a waterfall search interface, while verifying whether automated service dependency graphs and latency anomaly alerts are within project scope.
In the room: if asked to design Jaeger or Zipkin, verify early whether the interviewer prioritizes OpenTelemetry standard compliance, specific storage engines, and head vs tail sampling mechanics.
- Instrument services: Inject trace context including trace identifiers, span identifiers, and parent references into all inter-service network calls.
- Collect spans: Emit structured span records representing standard units of execution containing start times, durations, tags, and runtime metadata from every service.
- Correlate traces: Reconstruct the full request flow across multiple downstream microservices using a unified trace identifier.
- Visualize traces: Present waterfall and timeline visualizations of spans that display clear execution hierarchies and latency breakdowns.
- Search traces: Enable multi-attribute searching across trace identifiers, service names, operations, duration boundaries, error statuses, and custom tags.
- Service dependency graph: Discover and visualize topological service-to-service communication dependencies dynamically.
- Alerting: Trigger notifications on latency outliers and elevated error rates detected within the streaming trace telemetry path.
Non-Functional Requirements
Telemetry collection must never impede user-facing traffic. Call out the strict budget of less than 1% additional CPU utilization and less than 5% latency overhead on the application hot path early, and emphasize that intelligent sampling is an architectural requirement to control storage scale.
- Low Overhead: Tracing must introduce strictly less than 1% additional CPU utilization and less than 5% latency overhead to the critical execution path.
- High Throughput: The capacity baseline processes 5,000,000 raw spans per second system-wide, with the architecture scaling to the 10,000,000 spans per second staff stress case.
- Sampling Control: Support head-based, tail-based, and adaptive sampling strategies to balance diagnostic visibility against infrastructure cost.
- Scalability: Horizontally scale across thousands of microservices generating billions of telemetry events daily.
- Durability: Retain trace records for 14 days in the canonical planning model, with shorter hot-query tiers permitted when older data is retained more cheaply.
- Low Query Latency: Return trace search results in less than 1 second, and complete trace waterfall reconstructions in less than 2 seconds.
- Open Standards: Maintain full compatibility with OpenTelemetry data models and W3C Trace Context specifications.
Capacity Estimations
Telemetry volume grows rapidly in microservice architectures. Calculating raw vs sampled throughput demonstrates why sampling is a core system requirement rather than an optional optimization, because ingesting millions of spans per second without filtering creates unsustainable petabyte-scale storage costs.
| Metric | Calculation | Value |
|---|---|---|
| Microservices | Derived | 2,000 |
| Requests / sec (system-wide) | Given baseline workload | 500K |
| Avg spans per trace | Given (typical workload assumption) | 10 |
| Total spans / sec | 500K req/s x 10 spans/request | 5M (before sampling) |
| Effective retained trace rate | Target | 10% (1 in 10 traces) |
| Retained spans / sec | 5M x 10% retained-volume target | 500K |
| Avg span size | Given (typical workload assumption) | 500 bytes |
| Raw span ingress bandwidth | 5M raw spans/sec x 500 bytes | 2.5 GB/s |
| Retained ingestion throughput | 500K retained spans x 500 bytes | 250 MB/s |
| Active traces in tail window | 500K req/s x 30 to 60 seconds | 15 to 30M trace contexts before per-span and state overhead |
| Tail-sampling buffer (30 to 60s) | 5M raw spans/sec x 500 bytes x 30 to 60 seconds | ~75 to 150 GB raw buffer |
| Logical storage / day (before compression) | 250 MB/s x 86400 = 21.6 TB | 21.6 TB |
| Logical retention (14 days, before compression) | 21.6 TB x 14 ≈ 302 TB | ~300 TB |
Bandwidth and Processing Calculations: 1. Ingestion Volume (Sampled): - 500,000 spans/sec x 500 bytes/span = 250 MB/s. 2. Daily Storage Volume: - 250 MB/s x 86,400 seconds = 21.6 TB per day (raw JSON before compression). 3. Columnar ClickHouse Compression: - Highly compressed ClickHouse schemas reduce this by ~5–10x, requiring only ~3 TB/day.
Architecture Diagram
Microservices emit telemetry with shared trace context, transmitting spans through node-level agents into a centralized collector tier that batches and buffers the full-fidelity stream through Kafka before tail sampling, then indexes retained records into columnar storage and search backends for interactive visualization.
Component Deep Dives
Telemetry Pipeline Overview
Distributed tracing follows a span from context injection through agent aggregation, durable message buffering, tail-sampling decisions, and storage indexing. A lower-cost head-sampled mode is available when full-fidelity tail sampling is not required.
1. Trace Context Propagation
Distributed tracing relies on passing trace context across process boundaries. When an upstream service initiates an outbound call, the instrumentation SDK generates a globally unique trace_id and a local span_id, then injects metadata into the transport layer following the W3C Trace Context standard.
W3C traceparent structure:
traceparent: {version}-{trace-id}-{parent-span-id}-{trace-flags}
Example: 00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01
- version: 2 hex characters (currently 00).
- trace-id: 32 hex characters (globally unique per request flow).
- parent-span-id: 16 hex characters (identifies the caller's active span).
- trace-flags: 8-bit bitmap (01 indicates the trace is marked for sampling).Upon receiving the request, the downstream service extracts the traceparent header, preserves the trace_id, and instantiates a child span whose parent reference points directly to the caller's span_id.
2. Advanced Sampling Strategies ⭐
Ingesting 5M raw spans per second creates severe network bandwidth and storage overhead. To balance operational visibility with infrastructure cost, production systems implement a multi-layered hybrid sampling model.
| Strategy | How It Works | Pros | Cons |
|---|---|---|---|
| Head-based (probabilistic) | Decide at trace start (for example, sample 10% randomly) | Simple, low application overhead | May completely miss rare errors or low-volume endpoints |
| Tail-based ⭐ | Collect all spans for a trace, buffer them, and decide after the trace completes | Can retain 100% of eligible error or highly latent traces | Requires temporary in-memory collector buffering |
| Adaptive | Adjust the retention rate dynamically per service or endpoint based on traffic | Improves coverage of low-traffic endpoints | Requires more complex control-plane logic |
| Always-on for errors | Retain every trace that reaches the tail-sampling decision stage and registers an error | Preserves complete diagnostic coverage for eligible error traces | Error-heavy services increase storage footprint |
Optimal Production Mix: The canonical diagnostic path sends the full-fidelity span stream into the tail-sampling stage. It targets an average retained volume of about 10% while retaining eligible error traces and latency outliers. A separate 10% head-based mode remains available when lower overhead is more important than complete tail-sampling fidelity. Adaptive policies can increase retention during incidents.
3. Storage Engine Trade-offs
Selecting the appropriate persistence backend establishes the baseline for query latency, retention economics, and operational complexity across high-throughput telemetry workloads.
- Elasticsearch / OpenSearch: Indexes span attributes and tags into inverted indices, enabling flexible multi-attribute queries and rapid full-text searches. However, indexing overhead requires significant CPU and memory resources, driving high infrastructure expenses at scale.
- Cassandra: Delivers high write throughput with native time-to-live expiration per row. However, query capabilities are strictly limited to known primary keys, requiring queries to supply the exact
trace_idor maintain separate lookup tables. - ClickHouse: Employs columnar compression that typically achieves 5 to 10 times greater data compaction than Elasticsearch. It provides exceptional performance for analytical scans and duration aggregations, making it a primary choice for high-volume span persistence.
- Object Store with Compact Index (Grafana Tempo) ⭐: Persists immutable span blocks directly into object storage such as Amazon S3 or Google Cloud Storage, while maintaining lightweight Bloom filter or Redis indexes mapping
trace_idto storage chunk paths. This architecture reduces storage expenditures up to 100-fold and supports cost-effective long-term retention.
Event Bus Design (Kafka)
A partitioned event buffer decoupled through Kafka Architecture and Guarantees isolates collectors from ingestion backpressure, absorbs traffic bursts, and enables downstream workers to process raw and sampled telemetry independently.
Topic: spans-raw
Partitions: 128 # Partitioned by trace_id so all spans for one trace reach the same tail-sampling consumer partition
Partition key: trace_id
Retention: 24 hours # Durable buffer before the tail-sampling and storage stages
Replication factor: 3
Min insync replicas: 2
Topic: spans-sampled # Output after the tail-sampling decision
Partitions: 64
Partition key: trace_id # Preserves complete trace locality for downstream reconstruction
Retention: 14 days
Topic: span-metrics # Low-volume aggregates derived before trace sampling
Partitions: 64
Partition key: service_name
Retention: 24 hours
Producers: Central Collector metrics processor
Contains: request counts, error counts, latency histogram buckets, and parent-to-child edge counts
Producer:
Source: OpenTelemetry Collector exporting batched spans via OTLP to Kafka
Span Schema: "{ trace_id, span_id, parent_span_id, service_name, operation, duration_us, status_code, tags }"
Consumer Groups:
tail-sampler: Buffers each trace for 30 to 60 seconds, retains all eligible errors and latency outliers, and emits retained spans to spans-sampled
trace-writer: Writes retained spans from spans-sampled to Cassandra by trace_id for bounded trace reconstruction
span-writer: Writes retained spans from spans-sampled to ClickHouse in batches for analytical scans and retained-span storage
search-indexer: Indexes curated low-cardinality attributes from spans-sampled into Elasticsearch / OpenSearch for flexible search
dependency-graph: Consumes lightweight 100%-coverage parent-to-child edge summaries from span-metrics and updates the Redis service adjacency list
metrics-aggregator: Consumes lightweight 100%-coverage latency/error aggregates from span-metrics and uses Apache Flink to compute p50 and p99 metrics into ClickHouse
Ingestion Pipeline:
Path: "SDK to Node Agent to Central Collector to Kafka spans-raw to Tail Sampler to Kafka spans-sampled to Storage"
Head sampling: Optional 10% baseline mode when full-fidelity tail sampling is not required. Use it as a separate mode for the stream.
Tail sampling: Canonical production path receives the full spans-raw stream for error and latency-aware decisions
Dedupe key: (trace_id, span_id) so retried deliveries do not create duplicate spans
Dead Letter Queue: spans-raw-dlq after 3 delivery retries
Alerting: Trigger warning when consumer lag exceeds 60 secondsAPI Design
Tracing Domain Type Definitions
Strong typing across the instrumentation SDK, ingestion collector endpoints, and analytical query interfaces standardizes span serialization and trace assembly.
export interface SpanEventLog {
timestampUs: number;
message: string;
attributes?: Record<string, string | number | boolean>;
}
export interface SpanRecord {
traceId: string;
spanId: string;
parentSpanId?: string;
operationName: string;
serviceName: string;
serviceVersion?: string;
environment?: string;
startTimeUs: number;
durationUs: number;
statusCode: "OK" | "ERROR" | "UNSET";
tags: Record<string, string | number | boolean>;
logs?: SpanEventLog[];
process?: {
hostname?: string;
ip?: string;
};
}
export interface IngestSpansResponse {
accepted: number;
rejected: number;
errors?: string[];
}
export interface Trace {
traceId: string;
spans: SpanRecord[];
startTimeUs: number;
durationUs: number;
statusCode: "OK" | "ERROR" | "UNSET";
}
export interface TraceQueryResponse {
traces: Trace[];
nextPageToken?: string;
}
export interface TraceQueryFilter {
serviceName?: string;
operationName?: string;
minDurationUs?: number;
maxDurationUs?: number;
tags?: Record<string, string | number | boolean>;
startTimeUs: number;
endTimeUs: number;
limit?: number;
pageToken?: string;
}
export interface ServiceDependencyEdge {
parentService: string;
childService: string;
callCount: number;
errorCount: number;
p99LatencyMs: number;
}1. Ingest Spans (Collector API)
Services push batched spans over HTTP or gRPC using OpenTelemetry protocols to node agents or centralized collectors.
POST /api/v1/spans
Content-Type: application/json
[
{
"trace_id": "0af7651916cd43dd8448eb211c80319c",
"span_id": "b7ad6b7169203331",
"parent_span_id": "3a2fb4a1b3c4d5e6",
"operation_name": "GET /api/users",
"service_name": "user-service",
"start_time_us": 1710320000000000,
"duration_us": 12500,
"status_code": "OK",
"tags": { "http.method": "GET", "http.status_code": 200, "db.type": "postgresql" },
"logs": [{ "timestamp_us": 1710320000005000, "message": "cache miss, querying DB" }]
}
]
Response: 202 Accepted
{ "accepted": 1, "rejected": 0 }2. Query Traces
Operational dashboards query traces by service, operation, and duration bounds, or fetch complete span waterfalls by trace ID.
GET /api/v1/traces?service=user-service&operation=GET+/api/users&minDurationUs=100000&maxDurationUs=5000000&limit=20&pageToken=next-opaque-token&startTimeUs=1710320000000000&endTimeUs=1710406400000000
GET /api/v1/traces/0af7651916cd43dd8448eb211c80319c3. Service Dependency Discovery
Background streaming workers aggregate parent-to-child span links to build topological dependency graphs across microservices.
GET /api/v1/dependencies?startTimeUs=1710320000000000&endTimeUs=1710406400000000
Response: 200 OK
[
{ "parent": "api-gateway", "child": "user-service", "call_count": 150000, "error_count": 900, "p99_latency_ms": 180 },
{ "parent": "user-service", "child": "postgres-primary", "call_count": 120000, "error_count": 120, "p99_latency_ms": 75 }
]Common Error Responses
Standard error codes returned across the ingestion and query boundaries during rate limiting, malformed payloads, or query timeouts.
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
Span Document Schema (JSON Format)
Individual spans encapsulate operational metadata, microsecond timing precision, key-value tags, and chronological event logs.
{
"trace_id": "0af7651916cd43dd8448eb211c80319c",
"span_id": "b7ad6b7169203331",
"parent_span_id": "3a2fb4a1b3c4d5e6",
"operation_name": "GET /api/users/{id}",
"service_name": "user-service",
"service_version": "v2.3.1",
"environment": "production",
"start_time_us": 1710320000000000,
"duration_us": 12500,
"status_code": "OK",
"tags": {
"http.method": "GET",
"http.status_code": 200,
"http.url": "/api/users/123",
"db.type": "postgresql",
"db.statement": "SELECT * FROM users WHERE id = $1",
"peer.service": "postgres-primary",
"component": "spring-web"
},
"logs": [
{ "timestamp_us": 1710320000005000, "message": "cache_miss" },
{ "timestamp_us": 1710320000008000, "message": "db_query_complete", "attributes": { "rows": 1 } }
],
"process": {
"hostname": "user-svc-pod-abc123",
"ip": "10.0.5.42"
}
}Cassandra Schema Design
The primary table organizes spans by trace identifier, while secondary lookup tables partition service operations by start time and duration bucket for filtered range scans. Both tables use the same 14-day retention policy. Because the in-memory domain allows string, numeric, and boolean tag values, collectors normalize those values to strings when writing the Cassandra MAP and ClickHouse Map columns.
-- Trace lookup by trace_id
CREATE TABLE traces (
trace_id TEXT,
span_id TEXT,
parent_span_id TEXT,
operation_name TEXT,
service_name TEXT,
service_version TEXT,
environment TEXT,
start_time_us BIGINT,
duration_us BIGINT,
status_code TEXT,
tags MAP<TEXT, TEXT>,
logs_json TEXT,
process_json TEXT,
PRIMARY KEY (trace_id, span_id)
) WITH default_time_to_live = 1209600; -- Auto expire in 14 days
-- Secondary lookup table for service/operation/time/duration queries
CREATE TABLE service_operations (
service_name TEXT,
operation_name TEXT,
hour_bucket TEXT,
duration_bucket INT,
start_time_us BIGINT,
trace_id TEXT,
span_id TEXT,
duration_us BIGINT,
status_code TEXT,
PRIMARY KEY ((service_name, operation_name, hour_bucket, duration_bucket), start_time_us, trace_id, span_id)
) WITH CLUSTERING ORDER BY (start_time_us DESC) AND default_time_to_live = 1209600;ClickHouse Schema for Analytical Queries
ClickHouse stores retained spans in a time-partitioned columnar table for service, operation, duration, and latency analysis. The trace-oriented Cassandra table remains optimized for bounded trace reconstruction, while ClickHouse handles analytical scans over the retained span population.
CREATE TABLE spans (
trace_id String,
span_id String,
parent_span_id String,
operation_name LowCardinality(String),
service_name LowCardinality(String),
service_version LowCardinality(String),
environment LowCardinality(String),
start_time DateTime64(6),
duration_us UInt64,
status_code LowCardinality(String),
tags Map(String, String),
logs_json String,
process_json String
) ENGINE = MergeTree
PARTITION BY toStartOfHour(start_time)
ORDER BY (service_name, operation_name, start_time, trace_id, span_id)
TTL start_time + INTERVAL 14 DAY;Fault Tolerance
| Scenario Concern | System Solution Design |
|---|---|
| Collector Overload | OpenTelemetry agents implement backpressure, buffering up to 500 spans locally before applying the configured overload policy. The canonical full-fidelity path relies on the durable Kafka buffer and an explicit degraded mode when capacity is exhausted, preserving error traces first and shedding only the lowest-priority telemetry when the bounded buffer is full. |
| Storage Ingestion Failures | Stateless collector pools stream events to a highly durable, partitioned Kafka cluster before committing records to ClickHouse or object storage. |
| Application SDK Overhead | SDK operations execute asynchronously, and local agent communication uses bounded buffers to prevent blocking the hot execution path. |
| Unbounded Tag Cardinality | Centralized collectors strip dynamic user-specific tags at ingestion, redirecting high-cardinality values to unindexed log strings. |
1. Handling Clock Skew Across Services ⭐
Microservices run on virtual hosts with independent system clocks. A skew of 150ms could make a child span appear to start before its parent span began, breaking timeline visualizations.
- UI Adjustment: Use explicit parent-child relationships rather than wall-clock timestamps alone to render the execution tree sequentially.
- Relative Offsets: The SDK records span duration alongside explicit parent-child links so that the visualization layer orders elements by tree structure rather than relying solely on wall-clock readings. Telemetry export uses reliable gRPC or HTTP protocols with protocol buffers to stream spans to collectors instead of lossy datagrams.
- Clock Synchronization: Enforce network time protocol synchronization across cluster hosts to maintain clock skew under 10ms.
2. Trace Assembly Windows
Because spans arrive asynchronously and out of order across diverse microservices, the tail-sampling processor maintains an in-memory sliding assembly window of 30 to 60 seconds. When the root span reports its final duration, the collector waits for a bounded grace period for late child spans before committing the sampling decision. If the trace remains incomplete until the maximum 60-second window expires, the collector commits the available trace and marks it incomplete when later spans arrive. Traces that exceed the maximum assembly window are marked incomplete if later spans arrive after the decision.
Additional Considerations
Observer Effect Mitigation
To ensure the tracing SDK itself never degrades production request performance, the SDK employs a No-Op Span Pool. In the 10% head-sampled baseline mode, if the head-based sampler decides a request is not to be sampled, the SDK returns a statically allocated, lightweight no-op memory pointer. This avoids heap allocations and serialization overhead for approximately 90% of requests in that mode. The full-fidelity tail-sampling mode uses a different trade-off because spans must remain available for the later sampling decision.
Security and Access Control
Trace data can contain sensitive request metadata, so query APIs authenticate callers and authorize access by service, team, or tenant. Collector processors scrub credentials, authorization tokens, email addresses, and other sensitive fields before persistence. Collector and query traffic uses TLS, stored telemetry is encrypted at rest, and trace access is audited for compliance and incident review.
Interview Walkthrough
- 25-Minute Structured Interview Pace
Reserve deep storage tiering and cardinality management for staff-level inquiries.
- Scale Quantification (5 min): Establish that 500K requests per second generating 10 spans each yields 5M raw spans per second, proving that persistent storage of all spans is impractical without sampling.
- Tiered Architecture (6 min): Detail the three operational tiers, moving from client SDK context propagation through per-host agent batching to the centralized collector pipeline.
- Context Propagation (5 min): Explain W3C traceparent header injection and extraction across network boundaries so child spans reliably link into a unified execution tree.
- Hybrid Sampling Strategy (5 min): Explain the trade-off between a 10% head-based baseline mode and the full-fidelity tail-sampling path. Tail sampling is the mode that can retain all eligible errors and p99 latency outliers because it sees the complete trace.
- Resilient Storage and Buffering (4 min): Introduce Kafka between collectors and columnar storage engines like ClickHouse to decouple ingestion and absorb downstream maintenance spikes.
- Quantify Scale Upfront: Emphasize that 500,000 requests per second with an average of 10 spans per request generates 5,000,000 raw spans per second. Explain that retaining every span without sampling requires roughly 300 TB of logical storage over a 14-day retention window before compression and index overhead, which makes an intelligent sampling pipeline mandatory.
- Architect Three Clear Tiers: Structure the pipeline into application SDK context propagation, node-level agent batching, and centralized collector pools that publish to Kafka before tail-based sampling decisions. Call out that a 30 to 60 second tail window at the 500K request per second baseline represents 15 to 30 million active trace contexts before per-span and runtime overhead.
- Standardize Context Propagation: Highlight how injecting the W3C
traceparentheader across HTTP, gRPC, and message queue boundaries allows child spans to establish explicit parent links throughout the distributed service graph. - Justify Sampling Modes: Use a 10% probabilistic head sample when low overhead is the priority. When error fidelity and latency-aware retention are required, route the full-fidelity span stream through tail sampling so the decision stage can retain 100% of eligible error traces or p99 latency outliers.
- Decouple with Message Buffers: Place a partitioned message queue between collectors and storage backends using Kafka Architecture and Guarantees, ensuring that database maintenance or Elasticsearch re-indexing never drops in-flight spans.
- Evaluate Storage Trade-Offs: Contrast Elasticsearch for rich attribute search against columnar engines like ClickHouse and object storage solutions like Grafana Tempo, explaining how columnar compression achieves up to 100-fold cost reductions at hyperscale.
- Compensate for System Clock Skew: Explain why the visualization layer constructs execution timelines using explicit parent-child span relationships and duration measurements rather than relying solely on wall-clock timestamps.
- Address Common Interview Pitfalls: Avoid proposing 100% span retention without an explicit sampling model, because interviewers will immediately challenge the multi-petabyte storage cost and write amplification on the database tier.
Engineering Trade-offs
1. Why Introduce Kafka Durability Buffers?
Streaming spans directly from collectors into Elasticsearch or ClickHouse exposes storage clusters to thundering herd spikes during site-wide traffic surges. Introducing a partitioned Kafka topic between collectors and storage ingesters adds a modest 5 to 10 millisecond processing latency while providing robust buffer isolation. If storage clusters undergo maintenance or experience temporary index throttling, Kafka reliably buffers pending spans on disk, preventing data loss and smoothing ingestion traffic.
2. Head-Based vs Tail-Based Sampling Trade-offs
Head-based sampling is straightforward to implement and scales effortlessly because collectors require no stateful in-memory buffering. However, random head sampling decisions inevitably discard up to 90% of rare production errors and transient latency outliers. Conversely, tail-based sampling provides complete diagnostic fidelity for eligible error traces and latency outliers, but demands substantial in-memory buffer infrastructure across the collector tier to accumulate spans until the entire trace completes or the maximum decision window expires.
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.