Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

Logs Aggregation Patterns

A comprehensive guide to collecting, processing, storing, and querying logs at scale in distributed systems.


Overview

Log aggregation is the process of collecting logs from multiple sources, parsing and enriching them, storing them efficiently, and providing query capabilities. Modern distributed systems generate massive volumes of logs that require sophisticated aggregation patterns to extract actionable insights whilst managing costs.

Query & AnalysisStorage TiersProcessing PipelineCollection LayerLog SourcesApplicationsContainersInfrastructureNetwork DevicesAgents/ShippersSidecar ContainersSyslog ReceiversParsingEnrichmentFilteringRoutingHot StorageRecent/FrequentWarm StorageIndexed/CompressedCold StorageArchivedSearch UIDashboardsAlertsAnalyticsQuery & AnalysisStorage TiersProcessing PipelineCollection LayerLog SourcesApplicationsContainersInfrastructureNetwork DevicesAgents/ShippersSidecar ContainersSyslog ReceiversParsingEnrichmentFilteringRoutingHot StorageRecent/FrequentWarm StorageIndexed/CompressedCold StorageArchivedSearch UIDashboardsAlertsAnalytics

Structured Logging Conventions

Key Concepts

Structured logging outputs logs in a consistent, machine-parseable format (typically JSON) with standardised fields, making them easier to query, filter, and analyse.

Benefits:

  • Consistent field naming across services
  • Simplified parsing (no regex required)
  • Rich context without string concatenation
  • Efficient querying and indexing

Standard Log Schema

{
  "timestamp": "2024-12-07T10:30:00.123Z",
  "level": "ERROR",
  "service": "payment-service",
  "version": "v2.3.1",
  "environment": "production",
  "host": "pod-abc-123",
  "trace_id": "a1b2c3d4e5f6",
  "span_id": "7890xyz",
  "user_id": "user-12345",
  "message": "Payment processing failed",
  "error": {
    "type": "PaymentGatewayTimeout",
    "message": "Gateway did not respond within 5000ms",
    "stack_trace": "..."
  },
  "context": {
    "transaction_id": "txn-98765",
    "amount": 99.99,
    "currency": "GBP",
    "gateway": "stripe"
  },
  "duration_ms": 5234
}

Implementation Examples

Python (structlog):

import structlog

# Configure structured logger
structlog.configure(
    processors=[
        structlog.stdlib.filter_by_level,
        structlog.stdlib.add_logger_name,
        structlog.stdlib.add_log_level,
        structlog.stdlib.PositionalArgumentsFormatter(),
        structlog.processors.TimeStamper(fmt="iso"),
        structlog.processors.StackInfoRenderer(),
        structlog.processors.format_exc_info,
        structlog.processors.UnicodeDecoder(),
        structlog.processors.JSONRenderer()
    ],
    wrapper_class=structlog.stdlib.BoundLogger,
    context_class=dict,
    logger_factory=structlog.stdlib.LoggerFactory(),
    cache_logger_on_first_use=True,
)

log = structlog.get_logger()

# Structured logging with context
log.info(
    "payment_processed",
    transaction_id="txn-98765",
    amount=99.99,
    currency="GBP",
    gateway="stripe",
    duration_ms=234
)

Go (zap):

import (
    "go.uber.org/zap"
    "go.uber.org/zap/zapcore"
)

// Configure structured logger
config := zap.Config{
    Level:            zap.NewAtomicLevelAt(zap.InfoLevel),
    Encoding:         "json",
    OutputPaths:      []string{"stdout"},
    ErrorOutputPaths: []string{"stderr"},
    EncoderConfig: zapcore.EncoderConfig{
        TimeKey:        "timestamp",
        LevelKey:       "level",
        MessageKey:     "message",
        EncodeTime:     zapcore.ISO8601TimeEncoder,
        EncodeLevel:    zapcore.LowercaseLevelEncoder,
        EncodeDuration: zapcore.MillisDurationEncoder,
    },
}

logger, _ := config.Build()
defer logger.Sync()

// Structured logging
logger.Info("payment_processed",
    zap.String("transaction_id", "txn-98765"),
    zap.Float64("amount", 99.99),
    zap.String("currency", "GBP"),
    zap.String("gateway", "stripe"),
    zap.Int("duration_ms", 234),
)

Node.js (winston):

const winston = require('winston');

const logger = winston.createLogger({
  level: 'info',
  format: winston.format.combine(
    winston.format.timestamp({ format: 'YYYY-MM-DDTHH:mm:ss.SSSZ' }),
    winston.format.errors({ stack: true }),
    winston.format.json()
  ),
  defaultMeta: {
    service: 'payment-service',
    version: 'v2.3.1',
    environment: process.env.NODE_ENV
  },
  transports: [
    new winston.transports.Console()
  ]
});

// Structured logging
logger.info('payment_processed', {
  transaction_id: 'txn-98765',
  amount: 99.99,
  currency: 'GBP',
  gateway: 'stripe',
  duration_ms: 234
});

Log Levels

Use consistent severity levels across all services:

Level Usage Example
TRACE Very detailed debugging Function entry/exit, variable values
DEBUG Diagnostic information Configuration values, decision branches
INFO General informational Service started, request completed
WARN Potential issues Deprecated API usage, fallback triggered
ERROR Error conditions Failed operations, caught exceptions
FATAL Critical failures Service cannot continue, data corruption

Centralised Collection Pipelines

Collection Architecture

DestinationsAggregatorsShippersSourcesApp LogsSystem LogsContainer LogsFluentd/Fluent BitFilebeatVectorLogstashFluentd CentralVector AggregatorElasticsearchLokiS3/Object StorageDestinationsAggregatorsShippersSourcesApp LogsSystem LogsContainer LogsFluentd/Fluent BitFilebeatVectorLogstashFluentd CentralVector AggregatorElasticsearchLokiS3/Object Storage

Filebeat Configuration

Ship logs to Logstash:

# filebeat.yml
filebeat.inputs:
  - type: log
    enabled: true
    paths:
      - /var/log/app/*.log
    fields:
      service: payment-service
      environment: production
    fields_under_root: true

    # Multiline support for stack traces
    multiline.pattern: '^[0-9]{4}-[0-9]{2}-[0-9]{2}'
    multiline.negate: true
    multiline.match: after

  - type: container
    paths:
      - '/var/lib/docker/containers/*/*.log'
    processors:
      - add_kubernetes_metadata:
          host: ${NODE_NAME}
          matchers:
            - logs_path:
                logs_path: "/var/lib/docker/containers/"

# Output to Logstash
output.logstash:
  hosts: ["logstash:5044"]
  loadbalance: true
  compression_level: 3

# Buffer settings
queue.mem:
  events: 4096
  flush.min_events: 512
  flush.timeout: 5s

# Monitoring
monitoring.enabled: true

Fluentd/Fluent Bit Configuration

Fluent Bit as lightweight shipper:

# fluent-bit.conf
[SERVICE]
    Flush         5
    Daemon        off
    Log_Level     info
    Parsers_File  parsers.conf

[INPUT]
    Name              tail
    Path              /var/log/containers/*.log
    Parser            docker
    Tag               kube.*
    Refresh_Interval  5
    Mem_Buf_Limit     5MB
    Skip_Long_Lines   On

[FILTER]
    Name                kubernetes
    Match               kube.*
    Kube_URL            https://kubernetes.default.svc:443
    Kube_CA_File        /var/run/secrets/kubernetes.io/serviceaccount/ca.crt
    Kube_Token_File     /var/run/secrets/kubernetes.io/serviceaccount/token
    Merge_Log           On
    K8S-Logging.Parser  On
    K8S-Logging.Exclude On

[FILTER]
    Name    modify
    Match   *
    Add     cluster production-eu-west-1
    Add     environment production

[OUTPUT]
    Name            forward
    Match           *
    Host            fluentd-aggregator
    Port            24224
    Require_ack_response  true

Fluentd aggregator configuration:

# fluentd-aggregator.conf
<source>
  @type forward
  port 24224
  bind 0.0.0.0
</source>

# Parse JSON logs
<filter **>
  @type parser
  key_name log
  reserve_data true
  <parse>
    @type json
    time_key timestamp
    time_format %Y-%m-%dT%H:%M:%S.%NZ
  </parse>
</filter>

# Enrich with GeoIP
<filter **>
  @type geoip
  geoip_lookup_keys client_ip
  <record>
    location ${city.names.en["client_ip"]}
    coordinates ${location.latitude["client_ip"]},${location.longitude["client_ip"]}
  </record>
</filter>

# Route by log level
<match **>
  @type copy

  # Critical logs to PagerDuty
  <store>
    @type relabel
    @label @CRITICAL
  </store>

  # All logs to Elasticsearch
  <store>
    @type elasticsearch
    host elasticsearch
    port 9200
    index_name fluentd-${tag}-%Y%m%d
    type_name _doc
    <buffer tag, time>
      @type file
      path /var/log/fluentd-buffers/elasticsearch.buffer
      flush_mode interval
      flush_interval 5s
      flush_at_shutdown true
      retry_type exponential_backoff
      retry_timeout 1h
    </buffer>
  </store>
</match>

<label @CRITICAL>
  <filter **>
    @type grep
    <regexp>
      key level
      pattern /^(ERROR|FATAL)$/
    </regexp>
  </filter>

  <match **>
    @type http
    endpoint https://events.pagerduty.com/v2/enqueue
    <format>
      @type json
    </format>
  </match>
</label>

Vector Configuration

Modern, high-performance log aggregator:

# vector.toml
[sources.kubernetes_logs]
type = "kubernetes_logs"
auto_partial_merge = true

[transforms.parse_json]
type = "remap"
inputs = ["kubernetes_logs"]
source = '''
  . |= parse_json!(.message)
  .timestamp = to_timestamp!(.timestamp)
'''

[transforms.enrich]
type = "remap"
inputs = ["parse_json"]
source = '''
  .cluster = "production-eu-west-1"
  .environment = get_env_var!("ENVIRONMENT")

  # Add severity score
  .severity_score = if .level == "FATAL" {
    5
  } else if .level == "ERROR" {
    4
  } else if .level == "WARN" {
    3
  } else if .level == "INFO" {
    2
  } else {
    1
  }
'''

[transforms.sample_debug]
type = "sample"
inputs = ["enrich"]
rate = 10  # Only keep 10% of DEBUG logs

[transforms.filter_debug]
type = "filter"
inputs = ["sample_debug"]
condition = '.level == "DEBUG"'

[sinks.elasticsearch_hot]
type = "elasticsearch"
inputs = ["enrich"]
endpoint = "https://elasticsearch:9200"
mode = "bulk"
compression = "gzip"
bulk.index = "logs-%Y-%m-%d"

[sinks.s3_archive]
type = "aws_s3"
inputs = ["enrich"]
bucket = "logs-archive-bucket"
compression = "gzip"
encoding.codec = "json"
key_prefix = "logs/%Y/%m/%d/"
batch.max_bytes = 10485760  # 10MB

Parsing and Enrichment Strategies

Parsing Patterns

Logstash Grok Patterns:

# logstash.conf
filter {
  # Parse Apache access logs
  if [type] == "apache" {
    grok {
      match => {
        "message" => "%{COMBINEDAPACHELOG}"
      }
    }
    date {
      match => [ "timestamp", "dd/MMM/yyyy:HH:mm:ss Z" ]
    }
    geoip {
      source => "clientip"
    }
  }

  # Parse application logs (already JSON)
  if [type] == "application" {
    json {
      source => "message"
      target => "parsed"
    }
    mutate {
      rename => { "[parsed][timestamp]" => "@timestamp" }
      remove_field => [ "message" ]
    }
  }

  # Parse syslog
  if [type] == "syslog" {
    grok {
      match => {
        "message" => "%{SYSLOGTIMESTAMP:syslog_timestamp} %{SYSLOGHOST:syslog_hostname} %{DATA:syslog_program}(?:\[%{POSINT:syslog_pid}\])?: %{GREEDYDATA:syslog_message}"
      }
    }
  }

  # Custom application pattern
  grok {
    pattern_definitions => {
      "TRANSACTION_ID" => "[A-Z]{3}-[0-9]{8}"
    }
    match => {
      "message" => "Transaction %{TRANSACTION_ID:transaction_id} completed in %{NUMBER:duration:float}ms"
    }
  }
}

Enrichment Strategies

Raw LogParseAdd MetadataAdd ContextAdd Derived FieldsAdd GeoIPAdd Correlation IDsEnriched LogKubernetes LabelsCloud Provider TagsService VersionEnvironmentSeverity ScoreBusiness MetricsCountry/CityCoordinatesTrace IDSession IDRaw LogParseAdd MetadataAdd ContextAdd Derived FieldsAdd GeoIPAdd Correlation IDsEnriched LogKubernetes LabelsCloud Provider TagsService VersionEnvironmentSeverity ScoreBusiness MetricsCountry/CityCoordinatesTrace IDSession ID

Logstash Enrichment Pipeline:

filter {
  # Add Kubernetes metadata
  if [kubernetes][pod][name] {
    mutate {
      add_field => {
        "service" => "%{[kubernetes][labels][app]}"
        "namespace" => "%{[kubernetes][namespace]}"
        "pod" => "%{[kubernetes][pod][name]}"
        "container" => "%{[kubernetes][container][name]}"
      }
    }
  }

  # Lookup service details from cache/database
  elasticsearch {
    hosts => ["elasticsearch:9200"]
    index => "service-registry"
    query => "service:%{[service]}"
    fields => {
      "team" => "team"
      "owner" => "owner"
      "slack_channel" => "alert_channel"
    }
  }

  # Add derived fields
  ruby {
    code => '
      level = event.get("level")

      # Severity score for prioritisation
      severity_map = {
        "FATAL" => 5,
        "ERROR" => 4,
        "WARN"  => 3,
        "INFO"  => 2,
        "DEBUG" => 1
      }
      event.set("severity_score", severity_map[level] || 0)

      # Extract business metrics
      if event.get("message").include?("payment")
        event.set("category", "payment")
        event.set("business_critical", true)
      end
    '
  }

  # GeoIP enrichment
  if [client_ip] {
    geoip {
      source => "client_ip"
      target => "geo"
      fields => ["city_name", "country_name", "location"]
    }
  }

  # User agent parsing
  if [user_agent] {
    useragent {
      source => "user_agent"
      target => "ua"
    }
  }

  # Calculate request/response metrics
  if [duration_ms] {
    ruby {
      code => '
        duration = event.get("duration_ms").to_f

        # Classify latency
        if duration < 100
          event.set("latency_class", "fast")
        elsif duration < 1000
          event.set("latency_class", "normal")
        elsif duration < 5000
          event.set("latency_class", "slow")
        else
          event.set("latency_class", "critical")
        end
      '
    }
  }
}

Storage Retention and Tiering

Tiered Storage Strategy

Recent + HighPriorityAged LogsOld LogsCompliance7 days30 days90 daysSSD StorageFull IndexingFast QuerySSD/HDD MixCompressedIndexedHDD StorageHighly CompressedSearchableObject StorageCompressedImmutableIncoming LogsRoutingHot TierWarm TierCold TierArchive TierElasticsearch HotElasticsearch WarmElasticsearch FrozenS3 GlacierRecent + HighPriorityAged LogsOld LogsCompliance7 days30 days90 daysSSD StorageFull IndexingFast QuerySSD/HDD MixCompressedIndexedHDD StorageHighly CompressedSearchableObject StorageCompressedImmutableIncoming LogsRoutingHot TierWarm TierCold TierArchive TierElasticsearch HotElasticsearch WarmElasticsearch FrozenS3 Glacier

Elasticsearch Index Lifecycle Management (ILM)

{
  "policy": {
    "phases": {
      "hot": {
        "min_age": "0ms",
        "actions": {
          "rollover": {
            "max_size": "50GB",
            "max_age": "1d",
            "max_docs": 100000000
          },
          "set_priority": {
            "priority": 100
          }
        }
      },
      "warm": {
        "min_age": "7d",
        "actions": {
          "shrink": {
            "number_of_shards": 1
          },
          "forcemerge": {
            "max_num_segments": 1
          },
          "allocate": {
            "require": {
              "data": "warm"
            }
          },
          "set_priority": {
            "priority": 50
          }
        }
      },
      "cold": {
        "min_age": "30d",
        "actions": {
          "freeze": {},
          "allocate": {
            "require": {
              "data": "cold"
            }
          },
          "set_priority": {
            "priority": 0
          }
        }
      },
      "delete": {
        "min_age": "90d",
        "actions": {
          "delete": {}
        }
      }
    }
  }
}

Apply ILM policy to index template:

{
  "index_patterns": ["logs-*"],
  "template": {
    "settings": {
      "number_of_shards": 3,
      "number_of_replicas": 1,
      "index.lifecycle.name": "logs-policy",
      "index.lifecycle.rollover_alias": "logs-write",
      "codec": "best_compression"
    },
    "mappings": {
      "properties": {
        "@timestamp": { "type": "date" },
        "level": { "type": "keyword" },
        "service": { "type": "keyword" },
        "message": { "type": "text" },
        "trace_id": { "type": "keyword" },
        "duration_ms": { "type": "long" }
      }
    }
  }
}

Loki Retention Configuration

# loki-config.yaml
schema_config:
  configs:
    - from: 2024-01-01
      store: boltdb-shipper
      object_store: s3
      schema: v11
      index:
        prefix: loki_index_
        period: 24h

storage_config:
  boltdb_shipper:
    active_index_directory: /loki/index
    cache_location: /loki/cache
    shared_store: s3

  aws:
    s3: s3://eu-west-1/loki-chunks
    s3forcepathstyle: true

# Retention and compaction
compactor:
  working_directory: /loki/compactor
  shared_store: s3
  compaction_interval: 10m
  retention_enabled: true
  retention_delete_delay: 2h
  retention_delete_worker_count: 150

limits_config:
  # Global retention - can be overridden per tenant
  retention_period: 744h  # 31 days

  # Per-stream limits
  max_streams_per_user: 10000
  max_global_streams_per_user: 100000

  # Query limits
  max_query_length: 721h
  max_query_parallelism: 32
  max_entries_limit_per_query: 10000

# Per-tenant retention overrides
overrides:
  production:
    retention_period: 2160h  # 90 days
  staging:
    retention_period: 168h   # 7 days
  development:
    retention_period: 24h    # 1 day

S3 Lifecycle Policies for Log Archives

{
  "Rules": [
    {
      "Id": "TransitionToIA",
      "Status": "Enabled",
      "Prefix": "logs/",
      "Transitions": [
        {
          "Days": 30,
          "StorageClass": "STANDARD_IA"
        },
        {
          "Days": 90,
          "StorageClass": "GLACIER_IR"
        },
        {
          "Days": 180,
          "StorageClass": "DEEP_ARCHIVE"
        }
      ],
      "Expiration": {
        "Days": 2555
      },
      "NoncurrentVersionExpiration": {
        "NoncurrentDays": 1
      }
    }
  ]
}

Query Optimisation and Access Controls

Query Optimisation Strategies

Elasticsearch Query Best Practices:

{
  "query": {
    "bool": {
      "filter": [
        {
          "range": {
            "@timestamp": {
              "gte": "now-1h",
              "lte": "now"
            }
          }
        },
        {
          "term": {
            "service.keyword": "payment-service"
          }
        },
        {
          "terms": {
            "level.keyword": ["ERROR", "FATAL"]
          }
        }
      ]
    }
  },
  "aggs": {
    "errors_over_time": {
      "date_histogram": {
        "field": "@timestamp",
        "fixed_interval": "5m"
      },
      "aggs": {
        "error_types": {
          "terms": {
            "field": "error.type.keyword",
            "size": 10
          }
        }
      }
    }
  },
  "_source": ["@timestamp", "level", "message", "trace_id"],
  "size": 100,
  "sort": [
    { "@timestamp": "desc" }
  ]
}

Optimisation techniques:

// 1. Use filter context instead of query context (cacheable)
{
  "query": {
    "bool": {
      "filter": [  // Cached, no scoring
        { "term": { "status": 200 } }
      ],
      "must": [    // Scored, slower
        { "match": { "message": "error" } }
      ]
    }
  }
}

// 2. Limit fields returned
{
  "query": { ... },
  "_source": ["@timestamp", "message", "trace_id"],  // Only needed fields
  "docvalue_fields": ["duration_ms"]  // Use doc values for aggregations
}

// 3. Use index patterns for time-based queries
GET /logs-2024-12-*/_search  // Only search relevant indices

// 4. Pagination with search_after (not from/size)
{
  "query": { ... },
  "size": 100,
  "sort": [{ "@timestamp": "desc" }],
  "search_after": [1701936000000]  // Last result's sort value
}

Loki LogQL Optimisation:

# Optimised - filter early, then parse
{service="payment-service", level="error"}
  | json
  | transaction_id =~ "TXN-.*"
  | duration_ms > 5000
  | line_format "{{.timestamp}} {{.message}}"

# Avoid - parsing before filtering (slow)
{service="payment-service"}
  | json
  | level = "error"
  | transaction_id =~ "TXN-.*"

# Use metric queries for aggregations
sum by (service) (
  rate({environment="production"} |= "error" [5m])
)

# Label matchers are fastest - use them first
{cluster="prod", namespace="payments", pod=~"payment-.*"}

Index Optimisation

# Elasticsearch index settings for logs
PUT /logs-2024-12-07
{
  "settings": {
    "index": {
      "number_of_shards": 3,
      "number_of_replicas": 1,
      "refresh_interval": "30s",  // Reduce refresh for bulk indexing
      "codec": "best_compression",
      "translog.durability": "async",
      "translog.sync_interval": "30s"
    },
    "analysis": {
      "analyzer": {
        "log_analyzer": {
          "type": "custom",
          "tokenizer": "standard",
          "filter": ["lowercase", "stop"]
        }
      }
    }
  },
  "mappings": {
    "dynamic": "strict",  // Prevent mapping explosions
    "properties": {
      "@timestamp": {
        "type": "date",
        "format": "strict_date_optional_time||epoch_millis"
      },
      "message": {
        "type": "text",
        "analyzer": "log_analyzer",
        "norms": false  // Disable scoring for logs
      },
      "level": {
        "type": "keyword"
      },
      "service": {
        "type": "keyword"
      },
      "trace_id": {
        "type": "keyword",
        "index": true
      },
      "tags": {
        "type": "keyword"
      },
      "duration_ms": {
        "type": "long"
      },
      "metadata": {
        "type": "object",
        "enabled": false  // Don't index, just store
      }
    }
  }
}

Access Controls

Elasticsearch Role-Based Access Control:

{
  "roles": {
    "logs_reader": {
      "cluster": ["monitor"],
      "indices": [
        {
          "names": ["logs-*"],
          "privileges": ["read", "view_index_metadata"],
          "field_security": {
            "grant": ["@timestamp", "level", "service", "message"]
          },
          "query": {
            "term": {
              "environment": "production"
            }
          }
        }
      ]
    },
    "logs_admin": {
      "cluster": ["all"],
      "indices": [
        {
          "names": ["logs-*"],
          "privileges": ["all"]
        }
      ]
    },
    "team_payments_logs": {
      "cluster": ["monitor"],
      "indices": [
        {
          "names": ["logs-*"],
          "privileges": ["read"],
          "query": {
            "bool": {
              "should": [
                { "term": { "service": "payment-service" } },
                { "term": { "service": "billing-service" } },
                { "term": { "team": "payments" } }
              ]
            }
          }
        }
      ]
    }
  }
}

Loki Tenant Isolation:

# loki-config.yaml
auth_enabled: true

server:
  http_listen_port: 3100

limits_config:
  # Default limits
  ingestion_rate_mb: 10
  ingestion_burst_size_mb: 20

  # Per-tenant overrides
  per_tenant_override_config: /etc/loki/overrides.yaml

# overrides.yaml
overrides:
  team-payments:
    ingestion_rate_mb: 50
    ingestion_burst_size_mb: 100
    max_query_parallelism: 16

  team-platform:
    ingestion_rate_mb: 100
    ingestion_burst_size_mb: 200
    max_query_parallelism: 32

Nginx reverse proxy with authentication:

# nginx.conf
upstream elasticsearch {
    server elasticsearch:9200;
}

server {
    listen 443 ssl;
    server_name logs.example.com;

    ssl_certificate /etc/ssl/certs/logs.crt;
    ssl_certificate_key /etc/ssl/private/logs.key;

    # Authentication
    auth_request /auth;

    location /auth {
        internal;
        proxy_pass http://auth-service/verify;
        proxy_pass_request_body off;
        proxy_set_header Content-Length "";
        proxy_set_header X-Original-URI $request_uri;
    }

    location / {
        # Add tenant ID from auth
        proxy_set_header X-Tenant-ID $http_x_tenant_id;
        proxy_pass http://elasticsearch;

        # Rate limiting
        limit_req zone=logs_api burst=20 nodelay;

        # Only allow specific methods
        limit_except GET POST {
            deny all;
        }
    }

    # Block admin APIs
    location ~ ^/_cluster|_nodes|_cat {
        deny all;
        return 403;
    }
}

Quick Reference

Collection Tools Comparison

Tool Type Language Resource Usage Best For
Filebeat Shipper Go Low Simple file tailing, Elastic Stack
Fluent Bit Shipper C Very Low Kubernetes, embedded systems
Fluentd Aggregator Ruby Medium Complex routing, plugins
Logstash Aggregator Java High Rich filtering, Elastic Stack
Vector Both Rust Low High performance, observability
Promtail Shipper Go Low Loki-specific collection

Storage Solutions Comparison

Solution Type Query Language Strengths Weaknesses
Elasticsearch Full-text search DSL/SQL Powerful queries, rich ecosystem Resource intensive, complex
Loki Label-based LogQL Cost-effective, Grafana integration Limited full-text search
Splunk Enterprise SPL Feature-rich, enterprise support Expensive
CloudWatch Logs Cloud-native Insights Query AWS integration, managed Vendor lock-in, cost at scale
Datadog Logs SaaS Query syntax Easy setup, APM integration Expensive at scale

Common LogQL Queries

# Error logs in last hour
{service="payment-service"} |= "error" [1h]

# Rate of errors per minute
rate({level="error"}[5m])

# Count by service
sum by (service) (count_over_time({environment="production"}[1h]))

# Parse JSON and filter
{service="api"} | json | status_code >= 500

# Extract and aggregate metrics
sum by (endpoint) (
  avg_over_time({service="api"} | json | unwrap duration_ms [5m])
)

# Pattern matching
{app="nginx"} |~ ".*POST.*/api/.*"

# Multiple conditions
{cluster="prod"}
  | json
  | level="error"
  | duration_ms > 1000

Common Elasticsearch Queries

// Find errors in last hour
{
  "query": {
    "bool": {
      "filter": [
        {"range": {"@timestamp": {"gte": "now-1h"}}},
        {"term": {"level.keyword": "ERROR"}}
      ]
    }
  }
}

// Aggregate errors by service
{
  "size": 0,
  "aggs": {
    "services": {
      "terms": {"field": "service.keyword", "size": 20}
    }
  }
}

// Search with wildcard
{
  "query": {
    "wildcard": {
      "message": "*timeout*"
    }
  }
}

// Trace all logs for a request
{
  "query": {
    "term": {"trace_id.keyword": "abc123"}
  },
  "sort": [{"@timestamp": "asc"}]
}

Retention Recommendations

Environment Hot Tier Warm Tier Cold Tier Archive Total
Production 7 days 30 days 90 days 7 years 7+ years
Staging 3 days 7 days 14 days None 14 days
Development 1 day 3 days None None 3 days
Compliance 7 days 30 days 90 days 10 years 10+ years

Common Issues and Solutions

Issue: Log Ingestion Lag

Symptoms:

  • Delayed log visibility (minutes to hours)
  • Growing buffer sizes
  • Memory pressure on shippers

Solutions:

# Increase buffer and batch sizes
# filebeat.yml
queue.mem:
  events: 8192  # Increased from 4096
  flush.min_events: 2048
  flush.timeout: 3s  # Reduced from 5s

output.logstash:
  worker: 4  # Parallel workers
  bulk_max_size: 5120  # Larger batches
  compression_level: 3

# Logstash - increase workers
# logstash.yml
pipeline.workers: 16  # Match CPU cores
pipeline.batch.size: 500
pipeline.batch.delay: 50

# Add more Logstash instances (horizontal scaling)
# Use load balancer in Filebeat config

Issue: High Storage Costs

Solutions:

  1. Implement sampling for verbose logs:
# Logstash sampling
filter {
  if [level] == "DEBUG" {
    if [random_value] > 0.1 {  # Keep only 10% of DEBUG logs
      drop { }
    }
  }
}
  1. Use field exclusions:
{
  "mappings": {
    "_source": {
      "excludes": [
        "*.stack_trace",
        "metadata.large_field"
      ]
    }
  }
}
  1. Enable compression:
# Elasticsearch index template
{
  "settings": {
    "codec": "best_compression",
    "index.store.type": "hybridfs"
  }
}

Issue: Slow Query Performance

Symptoms:

  • Queries taking >10 seconds
  • Cluster CPU/memory spikes
  • Timeouts in dashboards

Solutions:

  1. Add appropriate indices:
# Create dedicated indices for keyword fields
PUT /logs-2024-12-07/_mapping
{
  "properties": {
    "trace_id": {
      "type": "keyword",
      "index": true,
      "doc_values": true
    }
  }
}
  1. Use index patterns instead of wildcards:
# Bad - searches all indices
GET /_all/_search

# Good - specific date range
GET /logs-2024-12-*/_search

# Better - exact indices
GET /logs-2024-12-06,logs-2024-12-07/_search
  1. Reduce cardinality:
# Logstash - normalize high-cardinality fields
filter {
  mutate {
    # Replace user IDs with hash
    gsub => ["message", "user_id=\d+", "user_id=REDACTED"]
  }
}

Issue: Missing or Incomplete Logs

Symptoms:

  • Gaps in log timelines
  • Missing stack traces
  • Truncated messages

Solutions:

  1. Configure multiline support:
# Filebeat multiline for Java stack traces
multiline.type: pattern
multiline.pattern: '^[[:space:]]+(at|\.{3})\b|^Caused by:'
multiline.negate: false
multiline.match: after
multiline.max_lines: 500
  1. Increase message size limits:
# Logstash
http.max_content_length: 200mb

# Fluentd
<source>
  @type forward
  chunk_size_limit 10m
</source>
  1. Check for back pressure:
# Monitor Filebeat registry
cat /var/lib/filebeat/registry/filebeat/log.json

# Check Logstash pipeline stats
curl -X GET "localhost:9600/_node/stats/pipelines?pretty"

Issue: Mapping Explosions

Symptoms:

  • Elasticsearch cluster becomes unresponsive
  • Index mapping has thousands of fields
  • Indexing failures

Solutions:

  1. Set strict mapping mode:
{
  "mappings": {
    "dynamic": "strict"
  }
}
  1. Use nested or flattened types:
{
  "mappings": {
    "properties": {
      "labels": {
        "type": "flattened"  // Prevents field explosion
      },
      "metadata": {
        "type": "object",
        "enabled": false  // Store but don't index
      }
    }
  }
}
  1. Configure field limits:
{
  "settings": {
    "index.mapping.total_fields.limit": 1000,
    "index.mapping.depth.limit": 20,
    "index.mapping.nested_fields.limit": 50
  }
}

Issue: Log Correlation Across Services

Solutions:

  1. Implement distributed tracing headers:
# Flask middleware
import uuid
from flask import request, g

@app.before_request
def before_request():
    g.trace_id = request.headers.get('X-Trace-ID', str(uuid.uuid4()))
    g.span_id = str(uuid.uuid4())

@app.after_request
def after_request(response):
    response.headers['X-Trace-ID'] = g.trace_id
    return response

# Always log trace_id
log.info("Processing request", trace_id=g.trace_id, span_id=g.span_id)
  1. Use consistent field names:
# Standardised field mapping across services
trace_id: UUID for request chain
span_id: UUID for this operation
parent_span_id: UUID of calling operation
service: Service name
operation: Operation name
  1. Create saved searches for trace correlation:
{
  "query": {
    "term": {
      "trace_id.keyword": "abc-123-def-456"
    }
  },
  "sort": [
    { "@timestamp": "asc" }
  ]
}

Related Topics

  • Observability Patterns - Metrics, traces, and logs integration
  • OpenTelemetry - Standardised instrumentation and collection
  • Prometheus - Metrics collection and alerting
  • Grafana - Visualisation and dashboards
  • Loki - Log aggregation system
  • Elasticsearch - Full-text search and analytics