深入 Logstash 12 - Multiple Pipelines 与 pipeline-to-pipeline
上一篇讲清了 DLQ 如何把逻辑上永久无法处理的 event 隔离到独立存储、等待补救管道重处理。这一篇进入另一个维度的隔离:一个 Logstash 进程里运行多条管道,管道之间互不干扰,又能通过 virtual address 串联成拓扑。
核心问题是三个:多管道的资源和生命周期如何隔离,pipeline-to-pipeline 通信怎么实现,以及什么时候该用多管道而不是在单管道里写 conditional。
单管道的局限
单管道(一个 logstash.conf)在简单场景下够用,但遇到两类需求时会暴露问题:
1 | |
Multiple Pipelines 的根本出发点是:不同数据流应该有各自独立的 queue、worker 池和生命周期,互不影响。
Multiple Pipelines 数据流
1 | |
每条 pipeline 有独立的 queue(含各自的 queue.type、queue.max_bytes)、独立的 worker 线程池(pipeline.workers)、独立的插件实例生命周期。一条管道崩溃或重载不影响其他管道。
pipelines.yml 配置
Multiple Pipelines 通过 config/pipelines.yml 定义,每个列表项对应一条 pipeline:
1 | |
pipeline.id 是 pipeline 的唯一标识,在 Node Stats API、DLQ 路径、pipeline-to-pipeline 地址里都会用到。path.config 可以指向单个文件,也可以用 glob 表达式指向一个目录(目录下所有 .conf 文件合并成一条 pipeline 的配置)。
当 pipelines.yml 存在时,logstash.yml 里的 path.config 和命令行 -e 参数不再生效——两种配置模式互斥。
pipeline-to-pipeline 通信
Multiple Pipelines 解决了隔离问题,但有时需要把多条 pipeline 串成拓扑:一条 pipeline 的 output 把 event 传给另一条 pipeline 的 input,整个过程在同一进程内完成,不经过 Kafka 或其他外部中间件。
这套机制叫 pipeline-to-pipeline,通过 virtual address 实现:
1 | |
send_to 是一个字符串列表,对应下游 pipeline input 里声明的 address。virtual address 只在进程内有效,不是网络地址,不需要端口。
pipeline output 把 event 推入下游 pipeline 的 input queue,下游的 queue 类型(memory / PQ)由下游 pipeline 自己的配置决定。
三种拓扑模式
1 | |
Distributor 是最常见的模式:一个统一入口承接所有 Beats/syslog 连接,避免每条子 pipeline 都开监听端口;再按 [type] 或其他元字段路由到对应的专门处理 pipeline。
最小实验:Distributor 拓扑
1 | |
1 | |
1 | |
router pipeline 开一个 Beats 监听端口,按 log_type 字段把 event 分发给 syslog-proc 或 nginx-proc。两条处理 pipeline 各自有独立的 worker 数和 queue 配置,互不干扰。
对应回内部对象
pipeline-to-pipeline 的 virtual address 机制在 Logstash 内部是一个进程内的有界阻塞队列(in-process queue),不是 loopback 网络连接。send_to 把 event 推入下游 pipeline 的 input 端 queue;如果下游 pipeline 配置了 PQ,event 进入 PQ 后具有 PQ 的 at-least-once 保证;如果下游是内存 queue,进程崩溃时 queue 里的 event 仍然会丢。
pipeline.workers 在 pipelines.yml 里针对每条 pipeline 单独配置,覆盖 logstash.yml 里的全局默认值。这是 Multiple Pipelines 最直接的资源隔离手段:高吞吐 pipeline 分配更多 worker,低频 pipeline 节省线程开销。
所有 pipeline 共享同一个 JVM 进程的 heap,-Xmx 设定的堆大小是全局的,不能按 pipeline 分配 heap 上限。内存 queue 的 queue.max_bytes 设定的是 queue 自身的字节上限,但 queue 里的 event 对象和 filter 处理时的临时对象都在共享 heap 上分配。pipeline 越多、heap 压力越大,这是多管道场景下 JVM 调优的核心约束。
模式提炼
1 | |
工程迁移表
| Logstash 概念 | Kafka 生态对应 | Flink 对应 | 通用 ETL 对应 |
|---|---|---|---|
| pipelines.yml(多管道定义) | consumer group 配置 / Kafka Streams topology | JobGraph 多算子链 | 多作业定义文件 |
| pipeline.id(管道标识) | consumer group id | JobVertex id | 作业名 / DAG 节点名 |
| pipeline.workers(每管道并发) | 消费者线程数 / partition 数 | 算子并行度 | 并发任务数 |
| pipeline-to-pipeline virtual address | 内部 topic / in-memory channel | 算子间 network buffer | 内存队列 / channel |
| Distributor 拓扑 | Router / filter by header | KeyedBroadcastProcessFunction 分流 | ETL 路由层 |
| Collector 拓扑 | 多 topic 同一 consumer group | union / connect 算子 | 多源合并作业 |
| 共享 JVM heap | 共享 broker / 共享 TaskManager | 共享 TaskManager heap | 共享执行引擎进程 |
Kafka Streams 的多拓扑(multiple topologies in one application)和 Logstash Multiple Pipelines 在设计意图上接近:在一个进程里并行运行多个处理逻辑,通过内部 channel 串联,避免跨网络的序列化开销。Flink 的算子链(operator chaining)在同一 TaskManager 线程内把多个算子合并成一个处理链,与 pipeline-to-pipeline 的"进程内直传"思路一致。核心取舍都是:进程内传递比网络传递快、但隔离性比独立进程弱,故障域扩展到整个进程。
常见误解
误解一:“Multiple Pipelines 等于多个 Logstash 进程”。Multiple Pipelines 是同一个 JVM 进程里的多条逻辑管道,共享 heap 和 OS 资源,不是多进程部署。真正的进程级隔离需要运行多个 Logstash 实例,代价是每个实例都有独立的 JVM overhead(内存、启动时间)。
误解二:“pipeline-to-pipeline 和写 Kafka 再读 Kafka 等价”。两者功能相似,但保证不同。pipeline-to-pipeline 的 virtual address 是进程内传递,进程崩溃时 in-process queue 里的 event 会丢(除非下游 pipeline 配置了 PQ);写 Kafka 再消费是跨进程的持久化传递,进程崩溃后 Kafka 里的数据不丢。需要跨进程可靠传递时,Kafka 是更合适的选择;纯粹的进程内路由和资源隔离时,pipeline-to-pipeline 更轻量。
误解三:“pipelines.yml 和 logstash.yml 可以同时指定 path.config”。两者互斥。当 pipelines.yml 文件存在且非空时,logstash.yml 里的 path.config 和命令行 -e 参数被忽略。混用会导致其中一种配置静默失效,排查时容易误判。
误解四:“多 pipeline 一定比单 pipeline 快”。Multiple Pipelines 的核心收益是隔离和可管理性,不是绝对性能提升。如果原本的单管道已经充分利用了 CPU 和 I/O,把它拆成多个 pipeline 只是把同样的线程换了个组织方式,吞吐不会提升,反而增加了 pipeline 间协调和调度的开销。选择多管道的理由应该是"这些数据流的处理要求差异大,需要独立的 worker/queue 配置",而不是"多管道=更快"。
练习
-
在本地用
pipelines.yml定义两条 pipeline:pipeline A 用内存 queue 跑generatorinput,pipeline B 用pipeline { address => "output-b" }input 接收来自 A 的 event,最后用stdout { codec => rubydebug }输出。观察 pipeline A 的 output 发出的 event 是否在 pipeline B 的 stdout 里出现,确认 virtual address 串联生效。 -
修改上面的实验,把 pipeline B 的
queue.type改成persisted,然后在 pipeline B 的 filter 里加一个sleep操作(用ruby { code => "sleep 0.5" })模拟处理慢。观察 pipeline A 的 generator 是否会被背压,还是继续发出 event 无视 B 的积压。解释 pipeline-to-pipeline 的背压传导边界。 -
思考题:给定一个场景——需要采集三个来源(syslog、nginx access log、应用 JSON log),分别做不同的 Grok 解析,但最终都写到同一个 ES 索引。用 Distributor 拓扑(一个 router pipeline + 三个 proc pipeline)和单管道(用 if/else if 分支)各自实现,列出两种方案在 worker 资源分配、queue 独立性、配置可维护性三个维度上的差异,决定哪种方案更适合这个场景。
系列导航
| 序号 | 主题 | 状态 |
|---|---|---|
| 00 | 导读:核心对象是 event,骨架是三段管道 | |
| 01 | Logstash 架构:JRuby、JVM 与 pipeline 的运行形态 | |
| 02 | event 模型:@timestamp、@metadata 与字段引用 |
|
| 03 | codec:字节流与 event 的边界转换 | |
| 04 | input 接入:文件、Beats、Kafka 与 TCP | |
| 05 | Grok:命名正则与 pattern 别名 | |
| 06 | dissect:分隔符切分与 Grok 的取舍 | |
| 07 | 常用 filter 组合:mutate、date、geoip、ruby | |
| 08 | Elasticsearch output:bulk 写入、索引模板与数据流 | |
| 09 | pipeline worker、batch 与背压:吞吐调优的三角 | |
| 10 | 内存队列 vs 持久队列:at-least-once 的边界 | |
| 11 | 死信队列(DLQ):无法处理的 event 去哪 | |
| 12 | Multiple Pipelines 与 pipeline-to-pipeline | 本篇 |
| 13 | 监控与诊断:Node Stats API 与瓶颈定位 | 下一篇 |
| 14 | 性能调优:JVM、worker、batch 的系统性方法 | |
| 15 | 演进一:Logstash vs Beats + Ingest Pipeline | |
| 16 | 演进二:Logstash vs Fluentd、Vector | |
| 17 | 演进三:Elastic Agent 与 Logstash 的分工 |
参考资料
- Logstash Multiple Pipelines 官方文档:https://www.elastic.co/guide/en/logstash/current/multiple-pipelines.html(pipelines.yml 结构、各字段语义、与 logstash.yml 的互斥关系)
- pipeline-to-pipeline 官方文档:https://www.elastic.co/guide/en/logstash/current/pipeline-to-pipeline.html(virtual address、三种拓扑模式、背压传导说明)
- pipeline input 插件文档:https://www.elastic.co/guide/en/logstash/current/plugins-inputs-pipeline.html(address 参数)
- pipeline output 插件文档:https://www.elastic.co/guide/en/logstash/current/plugins-outputs-pipeline.html(send_to 参数)
- Logstash 持久队列文档:https://www.elastic.co/guide/en/logstash/current/persistent-queues.html(per-pipeline queue 配置与 at-least-once 边界)
- Logstash 源码仓库:https://github.com/elastic/logstash(
logstash-core/lib/logstash/pipeline_bus.rb— pipeline-to-pipeline 内部实现) - 仓库内《深入 Elasticsearch》系列:Elasticsearch output 在多 pipeline 场景下的索引路由策略
