深入 Logstash 01 - 架构:JRuby、JVM 与 pipeline 的运行形态
.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 | |
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.options或LS_JAVA_OPTS。 - Warm-up 效应。JIT 编译在运行初期尚未生效,吞吐量在前几分钟通常低于稳态值。
插件体系:Ruby gem + logstash-plugin
插件以 Ruby gem 形式分发,gem 名称约定为 logstash-input-*、logstash-filter-*、logstash-output-*、logstash-codec-*。
1 | |
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 | |
每条 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 的队列总上限,默认 1024mb;queue.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 | |
- 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 | |
1 | |
启动日志中依次出现:
[logstash.runner]— JVM 参数打印[logstash.agent]— pipeline 注册[logstash.javapipeline]— pipeline 编译,.conf在这一步被编译成 JavaDataset执行图[logstash.filters.mutate]— 插件 register 回调[logstash.inputs.stdin]— input 线程就绪
在 stdin 输入任意文本后,rubydebug 输出显示完整 event 结构,包含 @timestamp、@version、host、message 和刚加上的 stage 字段。
pipeline 自身的运行时身份(id、worker 数、各阶段计数)不在 event 字段里,要看它得走 API:
1 | |
返回的 JSON 以 pipeline id 为 key(用 -f 启动时是 main),每个 id 下挂 events、plugins、reloads 三组计数。想确认 worker 数是否按命令行生效,读 /_node/pipelines?pretty 里的 workers 字段。
1 | |
关键对象映射
| 运行时概念 | 对应 Java/JRuby 类 | 说明 |
|---|---|---|
| event | org.logstash.Event |
Java 类,JRuby 侧通过 LogStash::Event 包装 |
| pipeline | org.logstash.execution.AbstractPipelineExt |
Ruby 侧对应 LogStash::JavaPipeline |
| 编译产物 | org.logstash.config.ir.CompiledPipeline |
.conf → PipelineIR → Dataset 执行图 |
| 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 也解决不了它的吞吐限制。
练习
-
用
--log.level=debug启动最小配置的 Logstash,量一次自己机器上的冷启动时间:取日志里[logstash.runner]第一行与Pipelines running那一行的时间差,同时记下 CPU 型号、bin/logstash-plugin list | wc -l的插件数和queue.type。再卸掉几个不用的插件重测一次,看插件数量对启动时间的影响有多大。 -
在
pipelines.yml中定义两条 pipeline,分别监听不同端口的 tcp input,用rubydebugcodec 输出,验证两条 pipeline 独立接收事件互不干扰;再用GET /_node/pipelines?pretty确认两条 pipeline 各自的workers、batch_size与 yml 里写的一致。 -
把
pipeline.workers显式设为 1,再改为 CPU 核数的两倍,用stdin快速粘贴大量带行号的文本,观察GET /_node/stats/pipelines里events.filtered的增长速率差异,同时对比两次输出里行号的顺序。然后在 workers 大于 1 的配置上加一行pipeline.ordered: true,确认 Logstash 拒绝启动并给出的错误信息。
系列导航
参考资料
- Logstash 执行模型文档:https://www.elastic.co/guide/en/logstash/current/pipeline.html(三段管道与 worker 的关系)
- logstash.yml 配置项清单:https://www.elastic.co/guide/en/logstash/current/logstash-settings-file.html(
pipeline.workers、pipeline.ordered、queue.*的取值与默认值) - Logstash 的 ECS 兼容模式:https://www.elastic.co/guide/en/logstash/current/ecs-ls.html(8.x 默认
v8,以及插件级/管道级/系统级三层覆盖顺序) - Logstash 多管道文档:https://www.elastic.co/guide/en/logstash/current/multiple-pipelines.html(
pipelines.yml可写哪些字段、未写的项如何回落到logstash.yml) - Logstash 持久队列文档:https://www.elastic.co/guide/en/logstash/current/persistent-queues.html(
path.queue、queue.max_bytes、queue.page_capacity的语义) - Logstash 版本发布说明:https://www.elastic.co/guide/en/logstash/current/releasenotes.html(Java 执行引擎从实验到默认、8.0 移除 Ruby 引擎的原文表述)
- Logstash 核心源码:https://github.com/elastic/logstash/tree/main/logstash-core(
lib/logstash/environment.rb是全部默认值的唯一出处,lib/logstash/java_pipeline.rb里的preserve_event_order?是pipeline.ordered的判定逻辑) - JRuby 调用 Java 的官方指南:https://github.com/jruby/jruby/wiki/CallingJavaFromJRuby(
Java::命名空间与类型映射规则)
