上一篇讲清了 DLQ 如何把逻辑上永久无法处理的 event 隔离到独立存储、等待补救管道重处理。这一篇进入另一个维度的隔离:一个 Logstash 进程里运行多条管道,管道之间互不干扰,又能通过 virtual address 串联成拓扑。

核心问题是三个:多管道的资源和生命周期如何隔离,pipeline-to-pipeline 通信怎么实现,以及什么时候该用多管道而不是在单管道里写 conditional。

单管道的局限

单管道(一个 logstash.conf)在简单场景下够用,但遇到两类需求时会暴露问题:

1
2
3
4
5
6
7
8
9
问题一:不同数据流的处理要求差异大
- syslog 流:高吞吐、无状态 filter、8 个 worker
- 慢查询日志:低频、复杂 Grok、2 个 worker
- 同一管道里 worker 数只能统一设,无法按流分配资源

问题二:一条 filter 处理慢导致全局堵塞
- 单管道共享一个 queue
- 某条 filter 插件慢(如重 geoip lookup),占用所有 worker slot
- 其他数据流的 event 在 queue 里积压,无法区分优先级

Multiple Pipelines 的根本出发点是:不同数据流应该有各自独立的 queue、worker 池和生命周期,互不影响。

Multiple Pipelines 数据流

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
                  Logstash 进程(单 JVM)
┌─────────────────────────────────────────────────────┐
│ │
│ Pipeline A │
│ ┌──────────┐ ┌──────┐ ┌────────────┐ │
│ │ input(s) │→ │queue │→ │filter+out │→ ES-A │
│ └──────────┘ └──────┘ └────────────┘ │
│ workers: 8 PQ: 512mb │
│ │
│ Pipeline B │
│ ┌──────────┐ ┌──────┐ ┌────────────┐ │
│ │ input(s) │→ │queue │→ │filter+out │→ ES-B │
│ └──────────┘ └──────┘ └────────────┘ │
│ workers: 2 memory queue │
│ │
└─────────────────────────────────────────────────────┘
共享:JVM heap、OS 文件句柄、CPU 调度
隔离:queue、worker 线程池、plugin 实例

每条 pipeline 有独立的 queue(含各自的 queue.typequeue.max_bytes)、独立的 worker 线程池(pipeline.workers)、独立的插件实例生命周期。一条管道崩溃或重载不影响其他管道。

pipelines.yml 配置

Multiple Pipelines 通过 config/pipelines.yml 定义,每个列表项对应一条 pipeline:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# config/pipelines.yml
- pipeline.id: syslog
path.config: "/etc/logstash/conf.d/syslog.conf"
pipeline.workers: 8
queue.type: persisted
queue.max_bytes: 512mb

- pipeline.id: slow-query
path.config: "/etc/logstash/conf.d/slow_query.conf"
pipeline.workers: 2
queue.type: memory

- pipeline.id: dlq-replay
path.config: "/etc/logstash/conf.d/dlq_replay.conf"
pipeline.workers: 1
queue.type: memory

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
2
3
4
5
6
7
8
9
10
11
12
13
# 上游 pipeline 的 output 段
output {
pipeline {
send_to => ["downstream-addr"]
}
}

# 下游 pipeline 的 input 段
input {
pipeline {
address => "downstream-addr"
}
}

send_to 是一个字符串列表,对应下游 pipeline input 里声明的 address。virtual address 只在进程内有效,不是网络地址,不需要端口。

pipeline output 把 event 推入下游 pipeline 的 input queue,下游的 queue 类型(memory / PQ)由下游 pipeline 自己的配置决定。

三种拓扑模式

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
拓扑一:Distributor(分发)
一条 input pipeline 接收所有数据,按条件分发到多条专门 pipeline

input-router pipeline:
input { beats { port => 5044 } }
output {
if [type] == "syslog" { pipeline { send_to => ["syslog-proc"] } }
else if [type] == "nginx" { pipeline { send_to => ["nginx-proc"] } }
else { pipeline { send_to => ["default-proc"] } }
}

拓扑二:Collector(汇聚)
多条数据源 pipeline 各自采集,汇聚到一条公共处理 pipeline

source-a pipeline: input → ... → output { pipeline { send_to => ["common"] } }
source-b pipeline: input → ... → output { pipeline { send_to => ["common"] } }
common pipeline: input { pipeline { address => "common" } } → filter → output

拓扑三:Forwarder(转发链)
pipeline A → pipeline B → pipeline C,逐级处理
每级 pipeline 只负责自己的处理逻辑,通过 address 串联

Distributor 是最常见的模式:一个统一入口承接所有 Beats/syslog 连接,避免每条子 pipeline 都开监听端口;再按 [type] 或其他元字段路由到对应的专门处理 pipeline。

最小实验:Distributor 拓扑

1
2
3
4
5
6
7
8
9
10
11
12
13
14
# pipelines.yml
- pipeline.id: router
path.config: "/etc/logstash/conf.d/router.conf"
pipeline.workers: 4

- pipeline.id: syslog-proc
path.config: "/etc/logstash/conf.d/syslog_proc.conf"
pipeline.workers: 6
queue.type: persisted
queue.max_bytes: 256mb

- pipeline.id: nginx-proc
path.config: "/etc/logstash/conf.d/nginx_proc.conf"
pipeline.workers: 2
1
2
3
4
5
6
7
8
9
10
11
12
# router.conf
input {
beats { port => 5044 }
}
filter {}
output {
if [@metadata][beat] == "filebeat" and [fields][log_type] == "syslog" {
pipeline { send_to => ["syslog-proc"] }
} else {
pipeline { send_to => ["nginx-proc"] }
}
}
1
2
3
4
5
6
7
8
9
10
11
# syslog_proc.conf
input {
pipeline { address => "syslog-proc" }
}
filter {
grok { match => { "message" => "%{SYSLOGLINE}" } }
date { match => ["timestamp", "MMM d HH:mm:ss", "MMM dd HH:mm:ss"] }
}
output {
elasticsearch { hosts => ["http://localhost:9200"] index => "syslog-%{+YYYY.MM.dd}" }
}

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.workerspipelines.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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
模式:进程内多租户管道 + virtual address 串联

隔离维度(每条 pipeline 独立):
- queue(类型、容量、落盘路径)
- worker 线程池(数量、batch size)
- plugin 实例(插件状态不跨 pipeline 共享)
- 生命周期(reload、崩溃不传染)

共享维度(全进程共享):
- JVM heap(-Xmx 全局上限)
- OS 资源(文件句柄、网络端口)
- CPU 调度(线程调度由 OS 决定,非绝对隔离)

串联原则:
- virtual address 是进程内 in-process queue,不是网络
- 下游 pipeline 的 queue 类型决定 event 在串联点的可靠性保证
- 拓扑复杂度随 pipeline 数量指数增长,超过 5-6 条时管理成本陡增

工程迁移表

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 配置",而不是"多管道=更快"。

练习

  1. 在本地用 pipelines.yml 定义两条 pipeline:pipeline A 用内存 queue 跑 generator input,pipeline B 用 pipeline { address => "output-b" } input 接收来自 A 的 event,最后用 stdout { codec => rubydebug } 输出。观察 pipeline A 的 output 发出的 event 是否在 pipeline B 的 stdout 里出现,确认 virtual address 串联生效。

  2. 修改上面的实验,把 pipeline B 的 queue.type 改成 persisted,然后在 pipeline B 的 filter 里加一个 sleep 操作(用 ruby { code => "sleep 0.5" })模拟处理慢。观察 pipeline A 的 generator 是否会被背压,还是继续发出 event 无视 B 的积压。解释 pipeline-to-pipeline 的背压传导边界。

  3. 思考题:给定一个场景——需要采集三个来源(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 的分工

参考资料