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.
graph TB
subgraph "Log Sources"
A1[Applications]
A2[Containers]
A3[Infrastructure]
A4[Network Devices]
end
subgraph "Collection Layer"
B1[Agents/Shippers]
B2[Sidecar Containers]
B3[Syslog Receivers]
end
subgraph "Processing Pipeline"
C1[Parsing]
C2[Enrichment]
C3[Filtering]
C4[Routing]
end
subgraph "Storage Tiers"
D1[Hot Storage<br/>Recent/Frequent]
D2[Warm Storage<br/>Indexed/Compressed]
D3[Cold Storage<br/>Archived]
end
subgraph "Query & Analysis"
E1[Search UI]
E2[Dashboards]
E3[Alerts]
E4[Analytics]
end
A1 & A2 & A3 & A4 --> B1 & B2 & B3
B1 & B2 & B3 --> C1
C1 --> C2 --> C3 --> C4
C4 --> D1
D1 --> D2 --> D3
D1 & D2 --> E1 & E2 & E3 & E4
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
flowchart LR
subgraph "Sources"
S1[App Logs]
S2[System Logs]
S3[Container Logs]
end
subgraph "Shippers"
F1[Fluentd/Fluent Bit]
F2[Filebeat]
F3[Vector]
end
subgraph "Aggregators"
L1[Logstash]
L2[Fluentd Central]
L3[Vector Aggregator]
end
subgraph "Destinations"
E1[Elasticsearch]
E2[Loki]
E3[S3/Object Storage]
end
S1 & S2 & S3 --> F1 & F2 & F3
F1 & F2 & F3 --> L1 & L2 & L3
L1 & L2 & L3 --> E1 & E2 & E3
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
flowchart TB
A[Raw Log] --> B{Parse}
B --> C[Add Metadata]
C --> D[Add Context]
D --> E[Add Derived Fields]
E --> F[Add GeoIP]
F --> G[Add Correlation IDs]
G --> H[Enriched Log]
C -.-> C1[Kubernetes Labels]
C -.-> C2[Cloud Provider Tags]
D -.-> D1[Service Version]
D -.-> D2[Environment]
E -.-> E1[Severity Score]
E -.-> E2[Business Metrics]
F -.-> F1[Country/City]
F -.-> F2[Coordinates]
G -.-> G1[Trace ID]
G -.-> G2[Session 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
flowchart LR
A[Incoming Logs] --> B{Routing}
B -->|Recent + High Priority| C[Hot Tier]
B -->|Aged Logs| D[Warm Tier]
B -->|Old Logs| E[Cold Tier]
B -->|Compliance| F[Archive Tier]
C -->|7 days| D
D -->|30 days| E
E -->|90 days| F
C -.->|SSD Storage<br/>Full Indexing<br/>Fast Query| C1[Elasticsearch Hot]
D -.->|SSD/HDD Mix<br/>Compressed<br/>Indexed| D1[Elasticsearch Warm]
E -.->|HDD Storage<br/>Highly Compressed<br/>Searchable| E1[Elasticsearch Frozen]
F -.->|Object Storage<br/>Compressed<br/>Immutable| F1[S3 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:
- Implement sampling for verbose logs:
# Logstash sampling
filter {
if [level] == "DEBUG" {
if [random_value] > 0.1 { # Keep only 10% of DEBUG logs
drop { }
}
}
}
- Use field exclusions:
{
"mappings": {
"_source": {
"excludes": [
"*.stack_trace",
"metadata.large_field"
]
}
}
}
- 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:
- 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
}
}
}
- 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
- 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:
- 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
- Increase message size limits:
# Logstash
http.max_content_length: 200mb
# Fluentd
<source>
@type forward
chunk_size_limit 10m
</source>
- 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:
- Set strict mapping mode:
{
"mappings": {
"dynamic": "strict"
}
}
- Use nested or flattened types:
{
"mappings": {
"properties": {
"labels": {
"type": "flattened" // Prevents field explosion
},
"metadata": {
"type": "object",
"enabled": false // Store but don't index
}
}
}
}
- 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:
- 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)
- 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
- 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