.conf 里声明的每个插件都是 Ruby gem,但把这些插件串起来执行的已经不是 Ruby 代码。Logstash 8.x 没有 Ruby 执行引擎——8.0 把它整个移除了。.conf 先由 ConfigCompiler 编译成 PipelineIR,再编译成由 Dataset 节点组成的 Java 执行图,插件的 filter/encode 方法是被这张图通过 JRuby 的 Java 互操作回调的。

上一篇确立了 event 与三段管道(input → filter → output)的核心抽象。这一篇顺着上面这道落差往下看两件事:插件是 Ruby gem 却跑在 JVM 上,在启动开销、线程模型和内存账上分别意味着什么;一个进程里多条 pipeline 的结构又是如何组织的。

版本前提

本系列的示例基于 Logstash 8.x,参数名与默认值以 8.19 分支的源码为准。

有一个设置会影响几乎每一篇的实验输出:pipeline.ecs_compatibility。Logstash 8 起它的默认值是 v8,所有实现了 ECS 兼容模式的插件都按 Elastic Common Schema 落位字段——stdin 的扁平 host 变成 [host][hostname],plain 与 line codec 额外写一个 [event][original],file input 的 path 变成 [log][file][path],geoip 的 country_code2 变成 [geo][country_iso_code]。要回到 8.0 之前的扁平字段名,在 logstash.yml 里设 pipeline.ecs_compatibility: disabled,或在单个插件上写 ecs_compatibility => disabled

本系列的实验一律按 disabled 模式给出输出,配置或命令行里会显式写出这一项,这样正文里的字段名和读者屏幕上的字段名能一一对上。凡是 ECS 落位差异会改变结论的地方,正文单独标注。

数据流全景

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
OS 进程边界
┌────────────────────────────────────────────────────────────────────┐
│ Logstash Process (JVM) │
│ │
│ ┌─ Pipeline A ───────────────┐ ┌─ Pipeline B ───────────────┐ │
│ │ input thread │ │ input thread │ │
│ │ │ │ │ │ │ │
│ │ ▼ │ │ ▼ │ │
│ │ queue(memory 或 PQ) │ │ queue(memory 或 PQ) │ │
│ │ │ batch │ │ │ batch │ │
│ │ ▼ │ │ ▼ │ │
│ │ worker 1 filter → output │ │ worker 1 filter → output │ │
│ │ worker 2 filter → output │ └────────────────────────────┘ │
│ │ worker N filter → output │ │
│ └────────────────────────────┘ │
│ │
│ JRuby runtime 单实例 · JVM heap 与 GC 由全进程共享 │
└────────────────────────────────────────────────────────────────────┘

PQ = Persistent Queue,queue.type: persisted 时启用;默认值是 memory,走内存队列。图里 worker 的标注写成 filter → output 而不是两组线程,这是 Logstash 线程模型里最容易记错的一点:pipeline.workers 控制的这一组线程同时执行 filter 和 output 两个阶段,进程里不存在独立的 output 线程池。

JRuby 与 JVM 的关系

Logstash 选择 JRuby 而非 CRuby/MRI,三条理由都落在 JVM 一侧:

  • Java 线程模型。CRuby 受 GIL(Global Interpreter Lock)限制,同一时刻只有一个线程执行 Ruby 字节码。JRuby 映射到真正的 OS 线程,多个 filter worker 线程可以并行执行插件代码。
  • Java 库互操作。logstash-core 的大量关键路径(codec 实现、queue、event 序列化)用 Java 写成,通过 JRuby 的 Java:: 命名空间直接调用,无需进程间通信。
  • JVM JIT。JRuby 代码经过 JVM 的 JIT 编译器优化,热路径与 Java 代码享有同等的即时编译待遇。

代价同样来自 JVM:

  • 启动延迟。JVM 初始化、JRuby 运行时加载、所有插件 gem 的 require 链依次串行完成,冷启动落在秒级到数十秒的区间,随已安装插件数量和 PQ 大小显著变化。官方文档没有给出启动耗时基准,要用的话自己测:记录日志里 [logstash.runner] 第一行到 Pipelines running 那一行的时间差,并且把机器型号、插件数量、queue.type 一起记下来,否则数字不可比。
  • GC 停顿。大批量事件在堆上分配时,Full GC 会造成处理延迟毛刺。Logstash 8.x 捆绑 JDK 17,config/jvm.options 里并没有显式指定收集器,实际用的是 JDK 默认的 G1;要换收集器或调堆大小,用 jvm.optionsLS_JAVA_OPTS
  • Warm-up 效应。JIT 编译在运行初期尚未生效,吞吐量在前几分钟通常低于稳态值。

插件体系:Ruby gem + logstash-plugin

插件以 Ruby gem 形式分发,gem 名称约定为 logstash-input-*logstash-filter-*logstash-output-*logstash-codec-*

1
2
3
4
5
6
7
8
# 查看已安装插件
bin/logstash-plugin list

# 安装新插件
bin/logstash-plugin install logstash-filter-translate

# 更新单个插件
bin/logstash-plugin update logstash-filter-mutate

logstash-plugin 本质上是封装过的 gem 命令,安装路径在 Logstash 自带的 vendor/bundle 目录内,与系统 Ruby 隔离。

插件类继承自 LogStash::Plugin,通过 config_name DSL 声明配置键,通过 register/run/filter/encode/decode 方法接入生命周期。logstash-core 在启动时反射加载这些类,拼装成 pipeline。

多 pipeline:pipelines.yml

自 Logstash 6.0 起,单个进程可以运行多条独立的 pipeline。配置文件 config/pipelines.yml 声明所有 pipeline:

1
2
3
4
5
6
7
8
9
10
11
- pipeline.id: main
path.config: "/etc/logstash/conf.d/main/*.conf"
pipeline.workers: 4
pipeline.batch.size: 125
pipeline.batch.delay: 50
queue.type: persisted

- pipeline.id: metrics
path.config: "/etc/logstash/conf.d/metrics/*.conf"
pipeline.workers: 2
queue.type: memory

每条 pipeline 持有独立的:

  • input 线程(或线程池)
  • 队列实例(内存或磁盘)
  • worker 线程池(数量由 pipeline.workers 控制,默认等于 CPU 核数;filter 与 output 两个阶段都跑在这组线程上)
  • output 插件实例(由 worker 线程调用,没有属于自己的线程池)

JRuby 运行时本身只有一个,所有 pipeline 共享同一个 JVM heap 和 GC。

在途 event 的内存账

单条 pipeline 在 filter 和 output 阶段同时在手的 event 数量有一个硬上界:pipeline.workers × pipeline.batch.size。默认值是 workers 等于 CPU 核数、batch.size 为 125,八核机器上约 1000 条。堆占用的粗算式就是把这个数乘以单事件在堆上的平均大小,而单事件大小不等于原始日志的字节数——grok 拆出的每个字段、tags 数组、@metadata 都各自占一份 Java 对象,通常是原始行长度的数倍。调大 workers 或 batch.size,这两个乘数一起放大堆压力。

这个乘积同时也是内存队列的容量来源。内存队列底层是一个有容量的 ArrayBlockingQueue,容量就等于 batch.size × workers,没有独立的大小旋钮可调。

换成持久队列之后,磁盘侧的边界由三个参数决定:path.queue 是队列目录,默认 <path.data>/queue(注意不是 queue.path);queue.max_bytes 是单条 pipeline 的队列总上限,默认 1024mbqueue.page_capacity 是单个 page 文件的大小,默认 64mb

无论用哪种队列,队列写满之后 input 线程的写入操作会阻塞,这个阻塞沿着协议一路传回上游:Beats 停止发送并重试,Kafka 停止拉取新 offset。Logstash 没有"丢弃在途 event 以保住吞吐"的降级路径,所以下游变慢必然表现为上游堵住,而不是数据静默消失。

pipeline.ordered:并行度的第一个代价

pipeline.ordered 控制 pipeline 是否保证输出顺序与输入顺序一致,默认值 auto,三种取值的行为差别很大:

  • auto:只有在显式pipeline.workers 设为 1 时才启用保序。单核机器上 workers 的默认值恰好是 1,但因为没有显式设置,auto 并不会启用保序——这个区分在源码里是 settings.set?("pipeline.workers") 那一次判断。
  • true:强制保序,并且要求单 worker。配了 true 而 workers 大于 1 时 Logstash 直接启动失败,报 enabling the 'pipeline.ordered' setting requires the use of a single pipeline worker
  • false:关掉保序相关的处理,省下这部分开销。

顺序保证与并行度在这里是互斥的,这也是第 03 篇"成帧不可逆"那条论证缺的另一半:一旦多 worker 并行取 batch,event 之间的先后关系在 filter 阶段就已经不可恢复,无论下游怎么补都补不回来。需要严格行序的场景只有两条路——单 worker,或者把顺序信息编进字段(序号、时间戳)后在查询侧排序。

pipeline 生命周期

1
2
3
loading      → running      → reloading (config 变化时)

stopped
  • loading:读取 .conf 文件,实例化插件,调用 register,分配 worker 线程。
  • running:input 线程持续产生 event,经队列进入 worker 线程,worker 执行 filter + output。
  • reloading:Logstash 支持 --config.reload.automatic,检测到配置文件变更时热重载,不重启进程;新 pipeline 实例替换旧实例,旧 worker 线程完成当前 batch 后退出。
  • stopped:收到 SIGTERM 后,input 停止接受新数据,worker 排干队列,output 刷新缓冲区,然后进程退出。

可运行实验

目标:观察 Logstash 启动序列和 pipeline 结构。

1
2
3
4
5
6
7
8
9
10
11
12
# 最小配置文件 /tmp/debug.conf
input {
stdin {}
}
filter {
mutate {
add_field => { "stage" => "filtered" }
}
}
output {
stdout { codec => rubydebug }
}
1
2
3
4
5
bin/logstash \
--log.level=debug \
--pipeline.ecs_compatibility=disabled \
-f /tmp/debug.conf \
--pipeline.workers=2

启动日志中依次出现:

  1. [logstash.runner] — JVM 参数打印
  2. [logstash.agent] — pipeline 注册
  3. [logstash.javapipeline] — pipeline 编译,.conf 在这一步被编译成 Java Dataset 执行图
  4. [logstash.filters.mutate] — 插件 register 回调
  5. [logstash.inputs.stdin] — input 线程就绪

在 stdin 输入任意文本后,rubydebug 输出显示完整 event 结构,包含 @timestamp@versionhostmessage 和刚加上的 stage 字段。

pipeline 自身的运行时身份(id、worker 数、各阶段计数)不在 event 字段里,要看它得走 API:

1
curl -s 'localhost:9600/_node/stats/pipelines?pretty'

返回的 JSON 以 pipeline id 为 key(用 -f 启动时是 main),每个 id 下挂 eventspluginsreloads 三组计数。想确认 worker 数是否按命令行生效,读 /_node/pipelines?pretty 里的 workers 字段。

1
2
# 查看已安装的插件 gem
bin/logstash-plugin list --verbose | head -30

关键对象映射

运行时概念 对应 Java/JRuby 类 说明
event org.logstash.Event Java 类,JRuby 侧通过 LogStash::Event 包装
pipeline org.logstash.execution.AbstractPipelineExt Ruby 侧对应 LogStash::JavaPipeline
编译产物 org.logstash.config.ir.CompiledPipeline .confPipelineIRDataset 执行图
worker thread java.lang.Thread JRuby 映射到真实 OS 线程
queue (memory) org.logstash.ext.JrubyWrappedSynchronousQueueExt 类名里没有 ackedqueue,也没有 ack 语义;底层是有容量的 ArrayBlockingQueue
queue (persisted) org.logstash.ackedqueue.Queue 磁盘 page + checkpoint,带 ack
plugin gem LogStash::Plugin 子类 Ruby 类,通过 JRuby Java 互操作调用 Java 核心

表里的每个 Java 类名都能在 elastic/logstash 仓库的 logstash-core/src/main/java/ 下按包路径打开对应文件,读源码时可以直接照着找。

模式提炼

Logstash 是"Ruby 接口 + Java 引擎"的混合架构。插件开发者用 Ruby DSL 声明配置和逻辑,底层执行由 Java 实现的 event、queue 和编译后的 Dataset 图驱动。最硬的一条证据是 8.0 那次移除:Ruby 执行引擎被整个删掉,而插件生态一个都没动——接口层和执行层的边界确实划在 DSL 与执行图之间。

多 pipeline 的独立性体现在:一条 pipeline 的 worker 阻塞不影响另一条;queue 故障只波及单条 pipeline;配置热重载可以针对单条 pipeline 执行。代价是共享 JVM heap,内存调优需要全局考虑。

工程迁移表

Logstash 概念 Kafka Streams 对应 Flink 对应 Fluentd 对应 Vector 对应
JRuby on JVM JVM 原生 Streams DSL JVM 原生 DataStream API CRuby on MRI(受 GIL 限制) Rust 原生二进制
plugin gem Processor/Transformer 接口 Function/ProcessFunction Fluent plugin gem component(Rust trait)
multi-pipeline 多个 KafkaStreams 实例 多个 Job 或 Slot Group multi-worker + tag 路由 多 topology
pipeline.workers stream threads task parallelism worker threads concurrency 参数
persistent queue Kafka topic 本身 State Backend (RocksDB) buffer plugin disk buffer

常见误解

误解一:“Logstash 就是个 Ruby 脚本,性能差是因为 Ruby 慢。”
实际情况:event 模型、队列和 pipeline 执行引擎都已经是 Java。Java 执行引擎的时间线是 6.1 以实验开关 --experimental-java-execution 引入、6.5 转 beta、7.0 成为默认、8.0 移除 Ruby 引擎。JRuby 层现在承担的是插件 DSL 解析和配置胶水,热路径由 Java 执行。

误解二:“多 pipeline 等于多进程,隔离彻底。”
实际情况:多 pipeline 在同一 JVM 进程内,共享 heap 和 GC。一条 pipeline 的内存泄漏会影响所有 pipeline。真正的进程隔离需要启动多个 Logstash 实例。

误解三:“pipeline.workers 设置越大越好。”
实际情况:调大 workers 的第一个代价是顺序,不是线程切换。pipeline.ordered: auto 下只要 workers 不是显式的 1,行序保证就没了。第二个代价是堆:在途 event 上界是 workers × batch.size,两个乘数任一放大都直接压在 heap 上。至于吞吐,官方对 workers 的说明恰恰是可以超过 CPU 核数,因为 output 常在 I/O wait 上空转;反过来说,当 output 本身就是瓶颈时,加 worker 也解决不了它的吞吐限制。

练习

  1. --log.level=debug 启动最小配置的 Logstash,量一次自己机器上的冷启动时间:取日志里 [logstash.runner] 第一行与 Pipelines running 那一行的时间差,同时记下 CPU 型号、bin/logstash-plugin list | wc -l 的插件数和 queue.type。再卸掉几个不用的插件重测一次,看插件数量对启动时间的影响有多大。

  2. pipelines.yml 中定义两条 pipeline,分别监听不同端口的 tcp input,用 rubydebug codec 输出,验证两条 pipeline 独立接收事件互不干扰;再用 GET /_node/pipelines?pretty 确认两条 pipeline 各自的 workersbatch_size 与 yml 里写的一致。

  3. pipeline.workers 显式设为 1,再改为 CPU 核数的两倍,用 stdin 快速粘贴大量带行号的文本,观察 GET /_node/stats/pipelinesevents.filtered 的增长速率差异,同时对比两次输出里行号的顺序。然后在 workers 大于 1 的配置上加一行 pipeline.ordered: true,确认 Logstash 拒绝启动并给出的错误信息。

系列导航

序号 主题
00 导读:核心对象是 event,骨架是三段管道
01 架构:JRuby、JVM 与 pipeline 的运行形态(本篇)
02 event 模型:@timestamp、@metadata 与字段引用
03 codec:字节流与 event 的边界转换
04 input 插件:拉取、监听与 Beats 接入
05 Grok 的本质:命名正则加预定义 pattern
06 dissect 与结构化 filter:放弃回溯换吞吐
07 常用 filter 组合:mutate、date、geoip 与条件
08 output 插件:Elasticsearch output 与批量写入
09 pipeline 执行模型:worker、batch 与背压
10 内存队列 vs 持久队列:可靠性的分界线
11 死信队列(DLQ):无法处理的 event 去哪
12 Multiple Pipelines 与 pipeline-to-pipeline
13 监控:Node Stats API、hot threads 与瓶颈定位
14 性能调优:JVM heap、批处理与持久队列磁盘
15 Logstash vs Beats vs Ingest Pipeline:该用谁
16 Logstash vs Fluentd vs Vector:日志管道的三种取舍
17 Logstash 的演进与 Elastic Agent 的冲击

参考资料