深入 Logstash 01 — 架构:JRuby、JVM 与 pipeline 的运行形态
上一篇确立了 event 与三段管道(input → filter → output)的核心抽象。这一篇进入 Logstash 的运行时层。运行时层容易被误解成"不就是个脚本引擎"。更准确的说法是:Logstash 是跑在 JVM 上的 JRuby 应用,插件是 Ruby gem,核心逻辑有相当比例用 Java 写成,两者通过 JRuby 的 Java 互操作机制在同一进程里协作。本文只抓一个问题:插件是 Ruby gem 却跑在 JVM 上意味着什么,以及一个 Logstash 进程里多 pipeline 的结构如何组织。
数据流全景
1 | |
PQ = Persistent Queue(可选,默认关闭时为内存队列)。
JRuby 与 JVM 的关系
Logstash 选择 JRuby(而非 CRuby/MRI)的原因不是 Ruby 语言本身,而是对 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 链,冷启动通常需要 30–60 秒。
- GC 停顿。大批量事件在堆上分配时,Full GC 会造成处理延迟毛刺;Logstash 8.x 默认使用 G1GC,可通过
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 核数) - output 实例
JRuby 运行时本身只有一个,所有 pipeline 共享同一个 JVM heap 和 GC。
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转成 Ruby 执行图)[logstash.filters.mutate]— 插件 register 回调[logstash.inputs.stdin]— input 线程就绪
在 stdin 输入任意文本后,rubydebug 输出显示完整 event 结构,包含 @timestamp、@version、host、message 字段。
1 | |
关键对象映射
| 运行时概念 | 对应 Java/JRuby 类 | 说明 |
|---|---|---|
| event | org.logstash.Event |
Java 类,JRuby 侧通过 LogStash::Event 包装 |
| pipeline | org.logstash.execution.JavaBasePipeline |
自 7.x 起核心逻辑已 Java 化 |
| worker thread | java.lang.Thread |
JRuby 映射到真实 OS 线程 |
| queue (memory) | org.logstash.ackedqueue.WrappedAckedQueue |
内存实现 |
| queue (persisted) | org.logstash.ackedqueue.AckedQueue |
磁盘 + WAL 实现 |
| plugin gem | LogStash::Plugin 子类 |
Ruby 类,通过 JRuby Java 互操作调用 Java 核心 |
模式提炼
Logstash 是"Ruby 接口 + Java 引擎"的混合架构。插件开发者用 Ruby DSL 声明配置和逻辑,底层执行由 Java 实现的 event、queue、pipeline 驱动。这种设计使插件生态保持低门槛(Ruby),同时关键路径享有 Java 的性能和多线程能力。
多 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 慢。”
实际情况:logstash-core 的 event 模型、队列、pipeline 执行引擎自 6.x 起已大量 Java 化。JRuby 层主要承担插件 DSL 解析和配置胶水,热路径由 Java 执行。
误解二:“多 pipeline 等于多进程,隔离彻底。”
实际情况:多 pipeline 在同一 JVM 进程内,共享 heap 和 GC。一条 pipeline 的内存泄漏会影响所有 pipeline。真正的进程隔离需要启动多个 Logstash 实例。
误解三:“pipeline.workers 设置越大越好。”
实际情况:worker 数超过 CPU 核数后,线程切换开销上升,而 I/O 密集型的 output 通常是瓶颈,增加 worker 不能解决 output 侧的吞吐限制。
误解四:“Logstash 重启很快,配置错了改完重启就行。”
实际情况:冷启动需要 30–60 秒(JVM + JRuby 初始化)。生产环境应使用 --config.reload.automatic 热重载,或在容器环境中做好预热策略。
练习
-
用
--log.level=debug启动最小配置的 Logstash,记录从进程启动到 input 就绪的完整日志行数和耗时;对比--log.level=info的输出量差异。 -
在
pipelines.yml中定义两条 pipeline,分别监听不同端口的 tcp input,用rubydebugcodec 输出,验证两条 pipeline 独立接收事件互不干扰。 -
修改
pipeline.workers为 1,再改为 CPU 核数的两倍,用stdin快速粘贴大量文本观察 Node Stats API(GET /_node/stats/pipelines)中events.filtered的增长速率差异。
系列导航
- 上一篇:深入 Logstash 00 — 导读:event、三段管道与学习路径
- 下一篇:深入 Logstash 02 — event 模型:@timestamp、@metadata 与字段引用
参考资料
- Elastic 官方文档 — How Logstash Works
- Elastic 官方文档 — Multiple Pipelines
- JRuby 官方文档 — JRuby and Java Integration
- Elastic 官方文档 — Logstash Plugin Development
- Elastic 博客 — Logstash Execution Engine
