ApiaryActive
Try: pause · settings · learn · wipe
← Community / Reading Room
LP
craft · 12 min read

Logstash Pipeline

Logstash is the beating heart of modern observability, turning raw, chaotic streams of data into structured, searchable events that power dashboards, alerts,…

Logstash is the beating heart of modern observability, turning raw, chaotic streams of data into structured, searchable events that power dashboards, alerts, and machine‑learning models. In a world where every sensor—from a server’s syslog to a beehive’s temperature probe—generates gigabytes of logs each day, the ability to parse and enrich those logs before they reach storage is not a luxury; it’s a prerequisite for any reliable analytics pipeline.

For the Apiary community, this matters twice over. First, researchers monitoring hive health need precise, time‑stamped metrics (temperature, humidity, colony weight) that are instantly comparable across apiaries worldwide. Second, the self‑governing AI agents that help coordinate conservation actions produce their own audit trails; without disciplined log handling, those agents can’t learn from their own decisions or be held accountable. A well‑crafted Logstash pipeline bridges the gap between raw telemetry and actionable insight, ensuring that every buzz and byte is captured with fidelity.

In this pillar article we’ll unpack the entire lifecycle of a Logstash pipeline—from the moment a line of text lands on a port, through parsing, enrichment, and finally into a durable store such as Elasticsearch. We’ll dive into concrete configurations, performance numbers, and real‑world examples that illustrate how Logstash can serve both traditional IT observability and the emerging needs of ecological data science.


1. What Is Logstash and Why It’s Central to the Elastic Stack

Logstash is an open‑source data collection engine written in JRuby that sits at the input → filter → output (IFO) junction of the Elastic Stack. According to the 2023 Elastic Stack Adoption Report, 68 % of Elastic customers run Logstash in production, handling an average of 2.4 TB of logs per month per deployment. Its primary responsibilities are:

RoleDescription
IngestionAccepts data from dozens of protocols (beats, syslog, HTTP, Kafka, JDBC, etc.).
ParsingTransforms free‑form text into structured fields using pattern matching, codecs, or custom scripts.
EnrichmentAdds contextual data (geo‑information, device metadata, threat intel) to each event.
RoutingSends events to one or more destinations: Elasticsearch, S3, Splunk, or even another Logstash instance.

Because Logstash runs on the JVM, it can leverage multi‑core CPUs and large heap sizes, making it suitable for high‑throughput environments (e.g., 100 k events/s on a 16‑core server with 32 GB heap). Its plugin ecosystem—over 200 official plugins as of version 8.12—covers virtually every data source and transformation need, from grok-patterns for syslog to geoip for IP‑to‑location mapping.


2. Core Architecture: Events, Pipelines, and Workers

At the heart of Logstash is the event object: a lightweight Ruby hash that carries a @timestamp, a message field (the raw payload), and any number of user‑defined fields. A pipeline is a directed graph that moves events through a series of filter plugins. Each pipeline runs in its own pipeline.id context, allowing multiple independent data flows on the same Logstash instance.

2.1 Pipeline Workers

Logstash spawns a pool of pipeline workers (pipeline.workers) that pull events from the input queue, apply filters, and push results to the output queue. By default this pool size equals the number of CPU cores. Benchmarks from Elastic’s own performance suite show:

CPU CoresWorkersThroughput (k events/s)Avg. Latency (ms)
444512
88929
16161807

Increasing workers beyond the core count yields diminishing returns due to GC pressure on the JVM heap.

2.2 Persistent Queues

When reliability is paramount—e.g., a beehive sensor network that may lose connectivity—Logstash can enable persistent queues (queue.type: persisted). Events are written to disk in a lock‑free segment file format, guaranteeing at‑least‑once delivery even if Logstash crashes. A typical configuration for a 10 GB queue on SSD can sustain ~250 k events before rotation, with a write latency of < 2 ms.

2.3 Multi‑Pipeline Mode

Since version 6.0, Logstash supports multiple pipelines defined in pipelines.yml. This enables separation of concerns: one pipeline may ingest raw IoT telemetry, while another handles security logs. Each pipeline can have its own pipeline.workers, queue.type, and even distinct plugin versions via pipeline-level settings.


3. Designing a Logstash Pipeline: From YAML to Runtime

A Logstash pipeline is defined in a .conf file using a straightforward DSL:

input {
  beats {
    port => 5044
  }
}
filter {
  grok {
    match => { "message" => "%{TIMESTAMP_ISO8601:timestamp} %{LOGLEVEL:level} %{GREEDYDATA:msg}" }
  }
  date {
    match => [ "timestamp", "ISO8601" ]
  }
}
output {
  elasticsearch {
    hosts => ["https://es01:9200"]
    index => "logs-%{+YYYY.MM.dd}"
    user => "logstash_system"
    password => "${ES_PASSWORD}"
  }
}

Key design considerations:

DecisionImpact
Input codec (e.g., json, plain, gzip)Determines how the raw payload is initially parsed.
Filter orderFilters are applied sequentially; placing heavy operations (e.g., ruby) later can reduce unnecessary work.
Conditional routing (if [type] == "beehive" )Prevents unrelated events from traversing expensive filters.
Pipeline-level settings (pipeline.batch.size, pipeline.batch.delay)Controls how many events are processed per batch, influencing throughput vs. latency.

A practical tip for large deployments: store common filter snippets in separate files and include them via the path option. This reduces duplication and eases updates across pipelines.


4. Parsing Techniques: Turning Text Into Structured Data

Parsing is the first transformation that converts a free‑form log line into searchable fields. Logstash offers several high‑performance parsers.

4.1 Grok – The Swiss‑Army Knife

Grok applies regular‑expression patterns stored in the patterns directory. The built‑in library includes over 1,200 patterns (e.g., COMMONAPACHELOG, NGINXERROR). A typical Apache access log parsing looks like:

filter {
  grok {
    match => { "message" => "%{COMBINEDAPACHELOG}" }
  }
}

Performance: In a 2024 benchmark on a 12‑core server, Grok parsed ~1.3 M events/s with a CPU utilization of 45 %. However, overly generic patterns can cause “catastrophic backtracking,” inflating latency. Use the named_captures_only => true flag to limit field creation.

4.2 Dissect – Zero‑Copy, High‑Speed Splitting

When the log format is delimiter‑based (CSV, pipe‑separated), dissect can be up to 5× faster than Grok because it avoids regex. Example for a beehive sensor CSV:

filter {
  dissect {
    mapping => {
      "message" => "%{timestamp},%{sensor_id},%{temperature},%{humidity},%{weight}"
    }
  }
  date {
    match => ["timestamp", "ISO8601"]
  }
}

Benchmarks show dissect handling 2.8 M events/s on a 16‑core box with negligible CPU load.

4.3 JSON and CSV Codecs

If the source already emits JSON (e.g., Beats or a REST API), the json codec on the input side eliminates the need for a filter:

input {
  http {
    codec => json
    port => 8080
  }
}

Similarly, the csv filter can parse delimited files with column headers, automatically casting numeric fields when convert => {"temperature" => "float"} is set.

4.4 Handling Unstructured Logs

When logs contain free‑form messages (e.g., application stack traces), a combination of mutate, split, and ruby filters can extract useful data. Example: extracting an exception class from a Java stack trace:

filter {
  if "Exception" in [message] {
    ruby {
      code => "
        m = event.get('message').match(/(?<exception>\\w+Exception):/)
        event.set('exception', m[:exception]) if m
      "
    }
  }
}

5. Enriching Data: Adding Contextual Value

Parsing yields fields, but enrichment adds meaning. Enrichment can be static (lookup tables) or dynamic (external APIs).

5.1 GeoIP – From IP to Location

The geoip filter uses MaxMind’s free GeoLite2 database (or a paid GeoIP2) to translate an IP address into latitude, longitude, city, and ASN. A typical security pipeline:

filter {
  geoip {
    source => "client_ip"
    target => "geo"
    database => "/usr/share/GeoIP/GeoLite2-City.mmdb"
  }
}

Numbers: In a 2023 security operation, enriching 5 M events per day with GeoIP added ≈ 0.8 ms per event, a negligible overhead given the value for geo‑visualization.

5.2 Translate – Static Key‑Value Lookups

The translate filter reads a CSV or YAML file to map codes to human‑readable values. For beehive telemetry, you might map sensor IDs to hive locations:

filter {
  translate {
    field => "sensor_id"
    destination => "hive_location"
    dictionary_path => "/etc/logstash/dictionaries/sensor_hive.csv"
    fallback => "unknown"
  }
}

A 10 k‑row dictionary loads into memory in < 30 ms and provides O(1) lookup time.

5.3 Enrich Plugin – External Data Sources

The enrich filter, introduced in Logstash 7.12, pre‑loads lookup data into an in‑memory cache, allowing join‑like operations without network latency. Example: joining a device inventory stored in Elasticsearch:

filter {
  enrich {
    policy => "device_inventory"
    field => "device_id"
    target => "device_meta"
    remove_fields => ["device_id"]
  }
}

The policy is built via the enrich API and refreshed every 15 minutes. In a high‑throughput IoT scenario (250 k events/s), the enrich filter added ~1.2 ms per event—still acceptable for batch analytics.

5.4 Ruby – Custom Logic

When built‑in filters fall short, the ruby filter offers full programmability. For AI agents that embed a trace‑id in logs, you can compute a deterministic hash:

filter {
  ruby {
    code => "
      require 'digest'
      trace = event.get('trace_id')
      event.set('trace_hash', Digest::SHA256.hexdigest(trace)) if trace
    "
  }
}

Performance tip: avoid heavy libraries inside the ruby block; pre‑require them in the ruby filter’s init => clause.


6. Managing State and Idempotency

Log ingestion pipelines often need to remember what they have already processed, especially when dealing with at‑least‑once delivery semantics.

6.1 Persistent Queues and Dead‑Letter Queues

A dead‑letter queue (DLQ) (dead_letter_queue.enable: true) captures events that repeatedly fail the output stage (e.g., Elasticsearch cluster outage). The DLQ writes events to a local dlq directory in JSON format, enabling later replay with the logstash-replay tool.

6.2 Checkpointing with Kafka Input

When Logstash reads from Kafka (kafka input plugin), it maintains offsets per topic/partition. By setting consumer_threads and auto_offset_reset => "earliest", you guarantee that no message is lost after a restart. The checkpoint file lives in data/queues/kafka and is persisted across restarts.

6.3 Idempotent Indexing

To avoid duplicate documents in Elasticsearch, use a document ID derived from a stable field (e.g., event_id). In the output:

output {
  elasticsearch {
    index => "beehive-%{+YYYY.MM.dd}"
    document_id => "%{event_id}"
    action => "create"
  }
}

If the same event arrives twice, Elasticsearch returns a 409 conflict, preventing duplication.


7. Scaling Logstash: From a Single Box to a Distributed Fleet

A single Logstash instance can handle hundreds of thousands of events per second, but production environments—especially those ingesting sensor data from thousands of hives—often require horizontal scaling.

7.1 CPU & Memory Sizing

Empirical data from the Elastic Performance Lab (2024) suggests the following baseline:

Target ThroughputRecommended CPUHeap Size
100 k events/s8‑core8 GB
250 k events/s16‑core16 GB
500 k events/s32‑core32 GB

Heap should be no more than 50 % of physical RAM to leave room for OS page cache.

7.2 Containerization and Kubernetes

Logstash runs natively in Docker; the official docker.elastic.co/logstash/logstash image includes all plugins. In Kubernetes, a StatefulSet is preferred when persistent queues are required, while a Deployment works for stateless pipelines.

Key manifest snippets:

apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: logstash
spec:
  serviceName: logstash
  replicas: 3
  selector:
    matchLabels:
      app: logstash
  template:
    metadata:
      labels:
        app: logstash
    spec:
      containers:
      - name: logstash
        image: docker.elastic.co/logstash/logstash:8.12.0
        env:
        - name: LS_JAVA_OPTS
          value: "-Xms16g -Xmx16g"
        volumeMounts:
        - name: pipelines
          mountPath: /usr/share/logstash/pipeline
        - name: queue
          mountPath: /usr/share/logstash/data
      volumes:
      - name: pipelines
        configMap:
          name: logstash-pipelines
      - name: queue
        persistentVolumeClaim:
          claimName: logstash-queue-pvc

Horizontal scaling is achieved by increasing replicas. Load balancing can be handled by a LoadBalancer Service in front of Beats or Kafka producers.

7.3 Load Balancing Input Sources

When using Beats, the beats input can accept connections from many agents. To avoid a single point of failure, deploy a HAProxy or NGINX TCP load balancer that forwards traffic to multiple Logstash pods. Metrics show that a 3‑node Logstash cluster with a round‑robin balancer can sustain ~600 k events/s with < 15 ms end‑to‑end latency.


8. Monitoring and Observability of Logstash Pipelines

A pipeline that is not observable is a risk. Logstash ships with a metrics API (http://localhost:9600/_node/stats) that exposes counters for input, filter, and output stages.

8.1 X‑Pack Monitoring

When X‑Pack is enabled, Logstash pushes its metrics to the Monitoring UI in Kibana. Important dashboards include:

  • Pipeline Throughput (events/s per pipeline)
  • Queue Size (current vs. max)
  • Filter Latency (average ms per filter plugin)

Alerts can be defined using Watcher or Alerting to trigger when queue.max_bytes exceeds 80 % of capacity.

8.2 Logstash Logs

Logstash writes its own logs to /var/log/logstash/logstash-plain.log. Use the log.level setting (INFO, DEBUG, WARN) to tune verbosity. For production, keep INFO and enable structured logging with json codec:

output {
  stdout { codec => json }
}

8.3 Prometheus Exporter

The community logstash-exporter scrapes the metrics endpoint and presents them to Prometheus. Sample query to alert on high filter latency:

avg_over_time(logstash_filter_processing_time_seconds_sum[5m]) / 
avg_over_time(logstash_filter_processing_time_seconds_count[5m]) > 0.05

9. Real‑World Use Cases

9.1 Bee‑Hive Sensor Data Pipeline

Apiary partners with a network of 2,300 smart hives across North America. Each hive streams a JSON payload every 30 seconds:

{
  "hive_id": "HB-1023",
  "timestamp": "2026-09-27T14:03:00Z",
  "temperature": 34.2,
  "humidity": 68,
  "weight": 42.5
}

Pipeline Overview

input {
  http {
    codec => json
    port => 8088
  }
}
filter {
  # Add hive location from static lookup
  translate {
    field => "hive_id"
    destination => "region"
    dictionary_path => "/etc/logstash/dictionaries/hive_region.csv"
  }
  # Convert temperature to Fahrenheit for legacy dashboards
  mutate {
    convert => { "temperature" => "float" }
    add_field => { "temp_f" => "%{temperature}" }
  }
  ruby {
    code => "event.set('temp_f', event.get('temperature') * 9/5 + 32)"
  }
}
output {
  elasticsearch {
    hosts => ["https://es-prod:9200"]
    index => "beehive-%{+YYYY.MM.dd}"
    user => "logstash_system"
    password => "${ES_PASSWORD}"
    pipeline => "beehive_enrich"
  }
}

Enrichment: The beehive_enrich pipeline adds weather data from an external API (via the http filter) and calculates a stress index (temp_f * humidity / weight). Researchers can now query “hives with stress index > 150” and receive real‑time alerts.

Performance: With 2,300 hives sending 2 events/minute, the pipeline processes ~77 k events/hour—well within a single‑node Logstash’s capacity, leaving headroom for future expansion.

9.2 AI Agent Audit Logging

Self‑governing AI agents in Apiary’s decision‑making layer emit logs like:

[2026-09-27 14:12:45.123][agent-07][TRACE] decision_id=5c9b9e3a action=deploy_sensor result=success latency=312ms

Parsing & Enrichment

filter {
  grok {
    match => { "message" => "\[%{TIMESTAMP_ISO8601:log_ts}\]\[%{DATA:agent_id}\]\[%{WORD:level}\] decision_id=%{DATA:decision_id} action=%{WORD:action} result=%{WORD:result} latency=%{NUMBER:latency}ms" }
  }
  date {
    match => ["log_ts", "ISO8601"]
    target => "@timestamp"
  }
  # Attach the current model version from a KV store
  translate {
    field => "agent_id"
    destination => "model_version"
    dictionary_path => "/etc/logstash/dictionaries/agent_model_version.csv"
  }
}

The enriched logs enable model‑drift analysis: correlating model_version with result over time highlights regressions, allowing the governance framework to trigger a rollback automatically.

9.3 Security Operations Center (SOC)

A typical SOC pipeline ingests firewall logs, DNS logs, and endpoint telemetry. Enrichment with threat intel (via the threatintel filter) adds a threat_score field. A sample rule in Kibana alerts when threat_score > 75 and geo.country_name == "China".


10. Best Practices and Common Pitfalls

PracticeReason
Keep filters statelessStateless filters (grok, dissect) scale linearly; stateful filters (ruby with external calls) can become bottlenecks.
Use conditionals earlyFilter only the
Frequently asked
What is Logstash Pipeline about?
Logstash is the beating heart of modern observability, turning raw, chaotic streams of data into structured, searchable events that power dashboards, alerts,…
What should you know about 1. What Is Logstash and Why It’s Central to the Elastic Stack?
Logstash is an open‑source data collection engine written in JRuby that sits at the input → filter → output (IFO) junction of the Elastic Stack. According to the 2023 Elastic Stack Adoption Report, 68 % of Elastic customers run Logstash in production, handling an average of 2.4 TB of logs per month per deployment.…
What should you know about 2. Core Architecture: Events, Pipelines, and Workers?
At the heart of Logstash is the event object: a lightweight Ruby hash that carries a @timestamp , a message field (the raw payload), and any number of user‑defined fields. A pipeline is a directed graph that moves events through a series of filter plugins . Each pipeline runs in its own pipeline.id context, allowing…
What should you know about 2.1 Pipeline Workers?
Logstash spawns a pool of pipeline workers ( pipeline.workers ) that pull events from the input queue, apply filters, and push results to the output queue. By default this pool size equals the number of CPU cores. Benchmarks from Elastic’s own performance suite show:
What should you know about 2.2 Persistent Queues?
When reliability is paramount—e.g., a beehive sensor network that may lose connectivity—Logstash can enable persistent queues ( queue.type: persisted ). Events are written to disk in a lock‑free segment file format, guaranteeing at‑least‑once delivery even if Logstash crashes. A typical configuration for a 10 GB…
References & sources
  1. Apiary Reading Room — Open, cited knowledge base — funded to keep bee & practical research free.
From the Apiary Reading Room. Opinion & editorial — not financial advice. We don't overclaim.
More from the Reading Room