•14 min read

OpenTelemetry Collector in Kubernetes: Scalable Architecture, Sampling & ClickHouse Export

OpenTelemetry Collector in Kubernetes: Scalable Architecture, Sampling & ClickHouse Export

This guide details the architectural considerations, configuration, and deployment strategies for the OpenTelemetry Collector within a Kubernetes environment, focusing on scalability, advanced sampling, and direct export to ClickHouse. The objective is to provide a robust, production-ready observability pipeline with zero data loss guarantees.

Audio Briefing
0:00 / 0:00

1. OpenTelemetry Collector Architectures in Kubernetes

Deploying the OpenTelemetry Collector in Kubernetes necessitates a clear architectural strategy. The primary models are Agent (DaemonSet), Gateway (Deployment), or a Hybrid approach. Each model presents distinct trade-offs in terms of resource utilization, latency, and processing capabilities.

1.1 Agent (DaemonSet) Architecture

The Agent architecture deploys an OpenTelemetry Collector instance on every Kubernetes node. This is typically achieved using a DaemonSet, ensuring a collector pod runs on each node. Agents are primarily responsible for local data collection and initial processing.

Characteristics:

  • Deployment Model: Kubernetes DaemonSet.
  • Data Collection: Collects telemetry data directly from applications running on the same node, often via host-level scraping or as a sidecar.
  • Processing: Performs lightweight, node-local processing (e.g., resource attribute enrichment, basic filtering, head-based sampling).
  • Export: Typically forwards data to a centralized Gateway Collector or directly to a backend.

Advantages:

  • Low Latency: Data travels a minimal path from application to collector.
  • Resilience: Node-level isolation; a failure in one agent does not impact other nodes.
  • Resource Isolation: Resources are distributed across nodes, preventing a single point of contention.
  • Network Efficiency: Reduces cross-node network traffic for initial collection.

Disadvantages:

  • Resource Overhead: Each node incurs the overhead of a collector instance.
  • Limited Global Context: Inefficient for advanced processing requiring a global view, such as tail-based sampling.
  • Management Complexity: Configuration changes require rolling updates across all nodes.

Use Cases:

  • Collecting host metrics and logs.
  • Collecting application traces with head-based sampling.
  • Initial data enrichment (e.g., adding Kubernetes metadata).

Example: OpenTelemetry Collector Agent DaemonSet

This example demonstrates a basic DaemonSet for an OTel Collector Agent, configured to receive OTLP and Prometheus metrics, and forward them.

apiVersion: opentelemetry.io/v1alpha1
kind: OpenTelemetryCollector
metadata:
  name: otel-agent
  namespace: observability
spec:
  mode: daemonset
  image: otel/opentelemetry-collector-contrib:0.90.1
  hostNetwork: true # Required for host-level scraping
  serviceAccount: otel-collector
  resources:
    requests:
      memory: "128Mi"
      cpu: "100m"
    limits:
      memory: "256Mi"
      cpu: "200m"
  config: |
    receivers:
      otlp:
        protocols:
          grpc:
          http:
      prometheus:
        config:
          scrape_configs:
            - job_name: 'kubernetes-nodes'
              static_configs:
                - targets: ['localhost:9100'] # Node Exporter
            - job_name: 'kubernetes-pods'
              kubernetes_sd_configs:
                - role: pod
              relabel_configs:
                - source_labels: [__meta_kubernetes_pod_annotation_prometheus_io_scrape]
                  action: keep
                  regex: true
                - source_labels: [__meta_kubernetes_pod_annotation_prometheus_io_path]
                  action: replace
                  target_label: __metrics_path__
                  regex: (.+)
                - source_labels: [__address__, __meta_kubernetes_pod_annotation_prometheus_io_port]
                  action: replace
                  regex: ([^:]+)(?::\d+)?;(\d+)
                  replacement: $1:$2
                  target_label: __address__
    processors:
      batch:
        send_batch_size: 10000
        timeout: 10s
      memory_limiter:
        limit_mib: 150
        spike_limit_mib: 50
        check_interval: 1s
      resourcedetection:
        detectors: [env, system, kubernetes]
        timeout: 2s
        override: true
    exporters:
      otlp:
        endpoint: "otel-gateway.observability.svc.cluster.local:4317" # Forward to Gateway
        tls:
          insecure: true
    service:
      pipelines:
        traces:
          receivers: [otlp]
          processors: [resourcedetection, batch, memory_limiter]
          exporters: [otlp]
        metrics:
          receivers: [otlp, prometheus]
          processors: [resourcedetection, batch, memory_limiter]
          exporters: [otlp]

1.2 Gateway (Deployment) Architecture

The Gateway architecture employs a centralized OpenTelemetry Collector deployment, typically as a Kubernetes Deployment, scaled horizontally. These collectors receive data from Agents or directly from applications, perform advanced processing, and export to various backends.

Characteristics:

  • Deployment Model: Kubernetes Deployment, often exposed via a Service.
  • Data Collection: Receives aggregated data from Agents or directly from applications.
  • Processing: Performs complex, resource-intensive operations (e.g., tail-based sampling, attribute manipulation, aggregation, fan-out).
  • Export: Exports processed data to long-term storage or analysis platforms (e.g., ClickHouse, Prometheus, Jaeger, Loki).

Advantages:

  • Centralized Processing: Ideal for operations requiring a global view, like tail-based sampling.
  • Resource Efficiency: Shared resources across multiple data streams, potentially reducing overall resource footprint compared to per-node agents for complex processing.
  • Simplified Management: Fewer instances to manage for core processing logic.
  • Scalability: Easily scaled horizontally by increasing replica count.

Disadvantages:

  • Increased Latency: Data must traverse the network twice (application -> agent -> gateway -> backend).
  • Single Point of Failure (if not scaled): A single gateway instance failure can disrupt the entire pipeline.
  • Network Overhead: All telemetry data flows through the gateway, potentially saturating network links.

Use Cases:

  • Tail-based sampling for traces.
  • Aggregating metrics from multiple sources.
  • Complex data transformation and filtering.
  • Exporting to multiple backend systems.

Example: OpenTelemetry Collector Gateway Deployment

This example shows a Deployment for an OTel Collector Gateway, configured to receive OTLP, process, and export.

apiVersion: opentelemetry.io/v1alpha1
kind: OpenTelemetryCollector
metadata:
  name: otel-gateway
  namespace: observability
spec:
  mode: deployment
  image: otel/opentelemetry-collector-contrib:0.90.1
  replicas: 3 # Scaled for high availability and throughput
  serviceAccount: otel-collector
  ports:
    - name: otlp-grpc
      port: 4317
      targetPort: 4317
      protocol: TCP
    - name: otlp-http
      port: 4318
      targetPort: 4318
      protocol: TCP
    - name: metrics
      port: 8888
      targetPort: 8888
      protocol: TCP
  resources:
    requests:
      memory: "1Gi"
      cpu: "500m"
    limits:
      memory: "2Gi"
      cpu: "1000m"
  config: |
    receivers:
      otlp:
        protocols:
          grpc:
          http:
    processors:
      batch:
        send_batch_size: 10000
        timeout: 10s
      memory_limiter:
        limit_mib: 1500
        spike_limit_mib: 500
        check_interval: 1s
      # Tail sampling configuration will be added here later
    exporters:
      # ClickHouse exporter will be added here later
      logging:
        verbosity: detailed
    service:
      pipelines:
        traces:
          receivers: [otlp]
          processors: [batch, memory_limiter] # Tail sampling will be inserted here
          exporters: [logging] # ClickHouse exporter will replace/augment this
        metrics:
          receivers: [otlp]
          processors: [batch, memory_limiter]
          exporters: [logging] # ClickHouse exporter will replace/augment this

The Hybrid architecture combines the strengths of both Agent and Gateway models. Agents collect data locally and forward it to Gateways, which then perform advanced processing and export. This is the recommended approach for most production environments.

Characteristics:

  • Agents (DaemonSet): Collect data, perform minimal processing (e.g., resource enrichment), and forward to Gateways.
  • Gateways (Deployment): Receive data from Agents, perform advanced processing (e.g., tail-based sampling, aggregation), and export to backends.

Advantages:

  • Optimal Performance: Low latency for initial collection, centralized powerful processing.
  • High Availability: Agents provide local resilience; Gateways can be scaled for high availability.
  • Scalability: Both layers can be scaled independently.
  • Flexibility: Allows for different processing logic at each stage.

Disadvantages:

  • Increased Complexity: Requires managing two distinct collector deployments and their interaction.
  • Higher Resource Footprint: Overall resource consumption is higher than a pure Gateway model due to per-node agents.

Example: Agent forwarding to Gateway (Agent config snippet)

The Agent configuration from Section 1.1 already demonstrates forwarding to a Gateway:

    exporters:
      otlp:
        endpoint: "otel-gateway.observability.svc.cluster.local:4317" # Forward to Gateway
        tls:
          insecure: true

1.4 Architectural Comparison Table

FeatureAgent (DaemonSet)Gateway (Deployment)Hybrid (Agent + Gateway)
Deployment ModelDaemonSet (per node)Deployment (centralized, scaled)DaemonSet (Agents) + Deployment (Gateways)
Resource UtilizationPer-node overhead, distributedCentralized, potentially higher per-instanceHigher overall, distributed collection, centralized processing
Latency (App->Collector)Negligible (local IPC/loopback)5-10ms (network hop)Negligible (App->Agent), then 5-10ms (Agent->Gateway)
Sampling CapabilityHead-based only, no global contextTail-based, probabilistic, global contextHead-based (Agent), Tail-based (Gateway)
ResilienceNode-level isolation, high local resilienceScaled for resilience, but central point of failure if not scaledHigh resilience at both layers
ComplexityLowMediumHigh
Best Use CaseHost metrics, logs, initial trace collectionAdvanced processing, global sampling, fan-outMost production environments, comprehensive observability
Data Loss RiskLower for local collection, higher for export if no bufferingHigher if not scaled or no bufferingLower with proper buffering at both layers
Advertisement

2. OpenTelemetry Collector Configuration Essentials

Effective OpenTelemetry Collector operation relies on a well-structured configuration. This section covers fundamental processors and exporters crucial for production deployments.

2.1 Core Components: Receivers, Processors, Exporters, Service

The OpenTelemetry Collector configuration is declarative and structured around four main components:

  • Receivers: How data gets into the collector (e.g., OTLP, Prometheus, Jaeger, Zipkin).
  • Processors: How data is transformed, filtered, or enriched within the collector (e.g., batch, memory_limiter, tail_sampling, resourcedetection).
  • Exporters: How data leaves the collector to various backends (e.g., OTLP, ClickHouse, Prometheus, Loki, Jaeger).
  • Service: Defines the pipelines, connecting receivers to processors and then to exporters for specific telemetry types (traces, metrics, logs).

2.2 Batching for Efficiency

The batch processor is critical for optimizing network utilization and reducing the load on backend systems. It aggregates telemetry data into larger batches before exporting.

Benefits:

  • Reduced Network Calls: Fewer, larger requests instead of many small ones.
  • Improved Throughput: Backends can process larger chunks of data more efficiently.
  • Lower CPU Usage: Less overhead per item for serialization and network I/O.

Configuration Parameters:

  • send_batch_size: The maximum number of telemetry items (spans, metric data points, log records) to send in a single batch.
  • timeout: The maximum duration to wait before sending a batch, even if send_batch_size is not reached.
  • send_batch_max_size: The absolute maximum size of a batch, overriding send_batch_size if necessary to prevent excessively large batches.

Example: batch processor configuration

    processors:
      batch:
        send_batch_size: 10000 # Max items per batch
        timeout: 10s         # Max wait time before sending
        send_batch_max_size: 12000 # Absolute max items, useful for high-volume bursts

2.3 Memory Limiter for Stability

The memory_limiter processor is essential for preventing the OpenTelemetry Collector from consuming excessive memory and being OOMKilled by Kubernetes. It monitors memory usage and, if limits are approached, drops data to prevent crashes.

Benefits:

  • Prevents OOMKills: Ensures collector stability under high load or misconfigurations.
  • Graceful Degradation: Prioritizes collector uptime over data completeness during memory pressure.

Configuration Parameters:

  • limit_mib: The maximum amount of memory (in MiB) the collector is allowed to use. When this limit is reached, data is dropped. This should be less than the container's Kubernetes memory limit.
  • spike_limit_mib: The maximum amount of memory (in MiB) the collector can temporarily exceed limit_mib by. This helps handle sudden spikes without immediate data drops.
  • check_interval: How frequently (duration) the memory usage is checked.

Example: memory_limiter processor configuration

    processors:
      memory_limiter:
        limit_mib: 1500 # Max memory in MiB before dropping data
        spike_limit_mib: 500 # Allowed spike above limit_mib
        check_interval: 1s # How often to check memory usage

Note: limit_mib should always be set lower than the Kubernetes container's limits.memory to allow the collector to react before Kubernetes OOMKills it. A good rule of thumb is limit_mib = limits.memory * 0.75.

2.4 Load Balancing Traces to Gateways

When using a Hybrid architecture, Agents need to distribute traces across multiple Gateway instances for high availability and load distribution. The loadbalancing exporter facilitates this.

Benefits:

  • High Availability: If one Gateway fails, agents automatically route to healthy ones.
  • Load Distribution: Spreads trace ingestion across multiple Gateway instances.
  • Scalability: Allows horizontal scaling of Gateways without reconfiguring agents.

Configuration Parameters:

  • protocol: The OTLP protocol to use (grpc or http).
  • resolver: Defines how target endpoints are discovered.
    • dns: Uses DNS SRV records or A records to discover endpoints. Recommended for Kubernetes services.
    • static: A fixed list of endpoints. Less dynamic.
  • routing_key: Determines how traces are routed. traceID is the default and recommended to ensure all spans of a single trace go to the same Gateway for proper processing (e.g., tail sampling).

Example: Agent loadbalancing exporter configuration

    exporters:
      otlp/loadbalancer:
        endpoint: "otel-gateway.observability.svc.cluster.local:4317" # Kubernetes Service DNS
        tls:
          insecure: true
        loadbalancing:
          resolver:
            dns:
              hostname: otel-gateway.observability.svc.cluster.local # Service DNS name
              port: 4317
          routing_key: traceID # Ensures all spans of a trace go to the same gateway
    service:
      pipelines:
        traces:
          receivers: [otlp]
          processors: [resourcedetection, batch, memory_limiter]
          exporters: [otlp/loadbalancer] # Use the loadbalancing exporter

3. Advanced Sampling Strategies: Tail-Based Sampling

Sampling is crucial for managing the volume of telemetry data, especially traces. While head-based sampling (deciding at the start of a trace) is simple, it often discards valuable context. Tail-based sampling addresses this by making sampling decisions after a trace has completed.

3.1 Why Tail-Based Sampling?

Tail-based sampling allows for intelligent sampling decisions based on the entire trace context. This means you can reliably capture:

  • Error Traces: All traces containing an error.
  • Slow Traces: All traces exceeding a certain latency threshold.
  • Specific Operations: Traces involving critical business operations.

Benefits:

  • Contextual Completeness: Ensures full traces are captured for analysis.
  • Targeted Data Collection: Focuses on traces that are most relevant for troubleshooting.

Trade-offs:

  • Centralized Processing: Requires a Gateway Collector to hold traces in memory until completion, increasing memory and CPU requirements on the Gateway.
  • Increased Latency: The decision is delayed until the trace is complete, which might add a small delay before the trace is exported.

3.2 tail_sampling Processor Configuration

The tail_sampling processor is configured within the Gateway Collector. It holds traces for a specified duration (decision_wait) to gather all spans before applying sampling policies.

Key Parameters:

  • decision_wait: The maximum duration to wait for all spans of a trace to arrive before making a sampling decision. This is critical for ensuring trace completeness.
  • num_traces: The maximum number of traces to hold in memory concurrently. This directly impacts memory usage.
  • expected_new_traces_per_sec: An estimate of the incoming trace rate, used for internal sizing.

Policy Types: The tail_sampling processor supports various policies, which can be combined using and or or composite policies:

  • always_sample: Always samples the trace.
  • drop_new: Drops new traces if num_traces is exceeded.
  • probabilistic: Samples traces based on a configured probability.
  • rate_limiting: Samples traces at a maximum rate per second.
  • status_code: Samples traces based on HTTP status codes (e.g., ERROR, UNSET).
  • latency: Samples traces exceeding a specified latency threshold.
  • attribute: Samples traces based on specific span or resource attributes.
  • composite: Combines multiple policies using and or or logic.

Example: tail_sampling processor with composite policy

This configuration samples all traces that contain an error OR have a latency greater than 500ms.

    processors:
      # ... other processors like batch, memory_limiter
      tail_sampling:
        decision_wait: 10s # Wait up to 10 seconds for all spans of a trace
        num_traces: 100000 # Max traces to hold in memory
        expected_new_traces_per_sec: 1000 # Expected incoming rate
        policies:
          - name: error-or-slow-trace
            type: composite
            composite:
              or_policies:
                - name: error-policy
                  type: status_code
                  status_code:
                    status_codes: [ERROR]
                - name: slow-trace-policy
                  type: latency
                  latency:
                    threshold_ms: 500 # Sample traces > 500ms
              on_failure: always_sample # If composite policy fails, always sample

Integration into Gateway Pipeline:

    service:
      pipelines:
        traces:
          receivers: [otlp]
          processors: [batch, memory_limiter, tail_sampling] # Insert tail_sampling here
          exporters: [logging] # Exporters will follow

4. Exporting to ClickHouse with Zero Data Loss

ClickHouse is an analytical database optimized for high-throughput ingestion and fast query execution, making it an excellent choice for storing OpenTelemetry traces and metrics. The clickhouse exporter allows direct integration.

4.1 ClickHouse Exporter Overview

The clickhouse exporter sends traces, metrics, and logs directly to ClickHouse tables. It supports various configurations for table schemas, TTLs, and compression.

4.2 Configuration for Traces

Traces are typically stored in a table designed to efficiently query by trace_id, service_name, operation_name, and time ranges.

Schema Considerations for ClickHouse Traces: A common schema for OpenTelemetry traces in ClickHouse includes:

  • trace_id (FixedString(16))
  • span_id (FixedString(8))
  • parent_span_id (FixedString(8))
  • service_name (String)
  • operation_name (String)
  • start_time_unix_nano (UInt64)
  • duration_nano (Int64)
  • kind (Int8)
  • status_code (Int8)
  • status_message (String)
  • attributes (Map(String, String)) - for span attributes
  • resource_attributes (Map(String, String)) - for resource attributes
  • events.timestamp_unix_nano (Array(UInt64))
  • events.name (Array(String))
  • events.attributes (Array(Map(String, String)))
  • links.trace_id (Array(FixedString(16)))
  • links.span_id (Array(FixedString(8)))
  • links.attributes (Array(Map(String, String)))

Example: ClickHouse exporter configuration for traces

    exporters:
      clickhouse/traces:
        dsn: "tcp://clickhouse-cluster.observability.svc.cluster.local:9000?database=otel"
        database: otel
        traces_table: otel_traces
        ttl: 7d # Data retention for 7 days
        compression: lz4 # Use LZ4 compression
        # Optional: Define specific schema mappings if default is not sufficient
        # trace_id_field: trace_id
        # span_id_field: span_id
        # ...

ClickHouse DDL for otel_traces table:

CREATE TABLE otel_traces (
    Timestamp DateTime64(9) CODEC(Delta, ZSTD(1)),
    TraceId FixedString(16) CODEC(ZSTD(1)),
    SpanId FixedString(8) CODEC(ZSTD(1)),
    ParentSpanId FixedString(8) CODEC(ZSTD(1)),
    TraceState String CODEC(ZSTD(1)),
    SpanName String CODEC(ZSTD(1)),
    SpanKind Int8 CODEC(ZSTD(1)),
    ServiceName LowCardinality(String) CODEC(ZSTD(1)),
    ResourceSchemaUrl String CODEC(ZSTD(1)),
    ResourceAttributes Map(LowCardinality(String), String) CODEC(ZSTD(1)),
    ScopeSchemaUrl String CODEC(ZSTD(1)),
    ScopeName String CODEC(ZSTD(1)),
    ScopeVersion String CODEC(ZSTD(1)),
    SpanAttributes Map(LowCardinality(String), String) CODEC(ZSTD(1)),
    DurationNanos UInt64 CODEC(ZSTD(1)),
    StatusCode Int8 CODEC(ZSTD(1)),
    StatusMessage String CODEC(ZSTD(1)),
    Events Nested (
        Timestamp DateTime64(9),
        Name String,
        Attributes Map(LowCardinality(String), String)
    ) CODEC(ZSTD(1)),
    Links Nested (
        TraceId FixedString(16),
        SpanId FixedString(8),
        TraceState String,
        Attributes Map(LowCardinality(String), String)
    ) CODEC(ZSTD(1))
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(Timestamp)
ORDER BY (ServiceName, SpanName, Timestamp, TraceId)
TTL Timestamp + INTERVAL 7 DAY
SETTINGS index_granularity = 8192, ttl_only_drop_parts = 1;

4.3 Configuration for Metrics

Metrics are typically stored in a table optimized for time-series

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement