深入 Logstash 12 - Multiple Pipelines 与 pipeline-to-pipeline
上一篇讲清了 DLQ 如何把逻辑上永久无法处理的 event 隔离到独立存储、等待补救管道重处理。这一篇进入另一个维度的隔离:一个 Logstash 进程里运行多条管道,管道之间互不干扰,又能通过 virtual address 串联成拓扑。
核心问题是三个:多管道的资源和生命周期如何隔离,pipeline-to-pipeline 通信怎么实现,以及什么时候该用多管道而不是在单管道里写 conditional。
单管道的局限
单管道(一个 logstash.conf)在简单场景下够用,但遇到两类需求时会暴露问题:
1 | |
Multiple Pipelines 的根本出发点是:不同数据流应该有各自独立的 queue、worker 池和生命周期,互不影响。第二类问题在官方文档里的措辞是 “Logstash, by default, is blocked when any single output is down”——这句话是多管道全部工程价值里最硬的一条动机。
Multiple Pipelines 数据流
1 | |
每条 pipeline 有独立的 queue(各自的 queue.type,PQ 还带各自的 queue.max_bytes)、独立的 worker 线程池(pipeline.workers)、独立的插件实例生命周期。PQ 和 DLQ 的存储位置按 pipeline.id 分命名空间,天然隔离。一条管道崩溃或重载不影响其他管道。
pipelines.yml 配置
Multiple Pipelines 通过 config/pipelines.yml 定义,每个列表项对应一条 pipeline:
1 | |
pipeline.id 是 pipeline 的唯一标识,在 Node Stats API、DLQ 路径、pipeline-to-pipeline 地址里都会用到。path.config 可以指向单个文件,也可以用 glob 表达式指向一个目录(目录下所有 .conf 文件合并成一条 pipeline 的配置)。
优先级的方向容易记反:不带参数启动时,Logstash 读 pipelines.yml 并实例化其中所有 pipeline;而一旦命令行带了 -e 或 -f,Logstash 就整个忽略 pipelines.yml,只在日志里留一条 warning。赢的一方是命令行参数。
两者也不是互斥关系。pipelines.yml 里某条 pipeline 没有显式写出的设置项,会回落到 logstash.yml 里的对应值——继承而非覆盖全部。上面例子里 dlq-replay 没写 pipeline.batch.size,它拿到的就是 logstash.yml(或内置默认)的 125。
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 自己的配置决定。
四种拓扑模式
官方给出的是四种,两两成对:distributor 与 collector 是数据的收敛方向相反,output isolator 与 forked path 是同一份数据扇出后下游做不做额外处理的区别。
1 | |
distributor 是最常见的模式:一个统一入口承接所有 Beats/syslog 连接,避免每条子 pipeline 都开监听端口;再按 [type] 或其他元字段路由到对应的专门处理 pipeline。
collector 是 distributor 的反向,代价也相反:配置被简化到一处,但所有数据挤回同一条 pipeline,隔离度随之消失。
forked path 的位置值得单独交代。在 pipeline input/output 出现之前,"一份 event 走两套处理规则"只能靠 clone filter 复制再用 if/else 分流,配置迅速变得难读;现在这种需求的标准写法就是扇出到多条下游 pipeline,各自写自己的 filter 段。
output isolator:多管道最硬的那个用例
output isolator 值得给完整配置,因为它是"问题二"的正面解法。一台机器同时把日志写 ES 和某个 HTTP 端点,HTTP 端点因为例行维护频繁不可用——单管道下这段时间 ES 也一起写不进去。
1 | |
1 | |
HTTP 端点不可用期间,event 堆在 buffered-http 自己的 PQ 里,buffered-es 照常写入。代价有三笔明账:磁盘占用最多翻倍(同一份数据进了两条 PQ)、序列化与反序列化成本约为单管道的三倍,以及一条兜底规则——任一下游 PQ 写满之后,两个 output 还是会一起停。PQ 买到的是缓冲时长,不是无限期解耦。
最小实验:Distributor 拓扑
1 | |
1 | |
1 | |
router pipeline 开一个 Beats 监听端口,按 log_type 字段把 event 分发给 syslog-proc 或 nginx-proc。两条处理 pipeline 各自有独立的 worker 数和 queue 配置,互不干扰。
投递语义与 ensure_delivery
串联点上到底什么时候阻塞、什么时候丢,由 pipeline output 的 ensure_delivery 参数和下游的状态共同决定。标准配置下 pipeline input/output 是 at-least-once,ensure_delivery 默认 true。
关键在于下游有两种"不通",官方对它们的处理完全不同。unavailable 指下游 pipeline 正在启动或正在 reload,只有这一种状态受 ensure_delivery 影响。blocked 指下游 pipeline 存在且在跑,但它内部某个插件卡住了,比如 output 正在等 ES 响应。
1 | |
这张表右下角那格是最容易踩的坑:把 ensure_delivery 设成 false 并不能让上游"不受下游影响"。它只覆盖 unavailable 一种情形,对 blocked 的下游无论取何值都照样阻塞上游。ensure_delivery => false 的真实用途是官方写明的那一个:想临时停掉某条下游 pipeline 而不牵连所有上游。
另外两条约束容易在 reload 时撞上。一条是避免成环:串联时要让数据保持单向流动,Logstash 关停时会等每条 pipeline 各自的工作完成,环形拓扑会让它永远等不到那一刻。另一条是 reload 的即时性:它会立刻按请求生效,包括删掉一条正在被上游写入的下游 pipeline。这会让上游阻塞,而且必须把下游恢复回来才能干净关停 Logstash。此时可以强制 kill,但在途 event 会丢,除非该 pipeline 开了 PQ。
对应回内部对象
virtual address 在 Logstash 内部是一个进程内的有界阻塞队列,实现在 PipelineBus / AbstractPipelineBus 两个 Java 类里。"有界"和"阻塞"这两点正是上面那张真值表的机制来源:队列满了写不进去,写入方就停在那里。
pipeline.workers 在 pipelines.yml 里针对每条 pipeline 单独配置,覆盖 logstash.yml 里的全局默认值。这是 Multiple Pipelines 最直接的资源隔离手段:高吞吐 pipeline 分配更多 worker,低频 pipeline 节省线程开销。
queue.max_bytes 是持久队列专属的:官方定义写的是"每条队列的总容量,以字节计",前提是 queue.type: persisted。内存队列没有字节记账这回事,它的在途上界由 pipeline.workers × pipeline.batch.size 决定,不能独立配置。
多管道的内存账不止 heap 一项。官方给了可以直接代入的整机估算:
1 | |
1 | |
对"多管道 + PQ"这个本篇主推的组合,第一块最容易被忽略:10 条 PQ pipeline 在 heap 之外还要多备出 1.4GB 上下的原生内存,而这部分没有任何 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 算子 | 多源合并作业 |
| output isolator 拓扑 | 每 sink 一个 topic + 独立 consumer | 每 sink 独立分支 + 各自 buffer | 每目标一条独立加载作业 |
| 共享 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 读 pipelines.yml;一旦带上 -e 或 -f,pipelines.yml 整个被忽略,只在日志里留一条 warning。用 -f 调试单条配置时以为多管道还在跑,就是这么来的。同时 pipelines.yml 与 logstash.yml 也不互斥:前者没写的设置项会回落到后者。
误解四:“ensure_delivery => false 能让上游不受下游影响”。它只覆盖 unavailable(下游启动中或 reload 中)一种情形。下游 pipeline 已经跑起来、只是内部插件卡住时属于 blocked,此时无论 ensure_delivery 取什么值都会阻塞上游。想真正把故障域切开,靠的是给下游各自配 PQ(output isolator 拓扑),而 PQ 也只是把阻塞推迟到队列写满。
误解五:“多 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 是否会被背压。接着把 pipeline A 的 output 改成pipeline { send_to => ["output-b"] ensure_delivery => false }重跑一次,确认背压依然存在——此时 B 属于 blocked 而非 unavailable。最后把 B 从pipelines.yml里删掉触发 reload,观察 A 是否阻塞,以此把真值表的两列各验一遍。 -
思考题:给定一个场景——需要采集三个来源(syslog、nginx access log、应用 JSON log),分别做不同的 Grok 解析,但最终都写到同一个 ES 索引。用 Distributor 拓扑(一个 router pipeline + 三个 proc pipeline)和单管道(用 if/else if 分支)各自实现,列出两种方案在 worker 资源分配、queue 独立性、配置可维护性三个维度上的差异,决定哪种方案更适合这个场景。
系列导航
参考资料
- 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、四种拓扑模式、投递保证与
ensure_delivery、avoid cycles) - 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、ensure_delivery 参数)
- Logstash 持久队列文档:https://www.elastic.co/guide/en/logstash/current/persistent-queues.html(per-pipeline queue 配置、
queue.max_bytes的 PQ 限定、at-least-once 边界) - Logstash JVM 设置文档:https://www.elastic.co/guide/en/logstash/current/jvm-settings.html(off-heap 构成与多管道 + PQ 的整机内存估算公式)
- Logstash 源码仓库:https://github.com/elastic/logstash(
logstash-core/src/main/java/org/logstash/plugins/pipeline/PipelineBus.java与同目录AbstractPipelineBus.java— pipeline-to-pipeline 的进程内队列实现) - 仓库内《深入 Elasticsearch》系列:Elasticsearch output 在多 pipeline 场景下的索引路由策略
