深入 Logstash 09 - pipeline 执行模型:worker、batch 与背压
上一篇(08)覆盖了 Elasticsearch output 的 bulk 写入机制。这一篇进入 pipeline 执行模型本身:pipeline.workers、pipeline.batch.size、pipeline.batch.delay 三个参数各自控制什么,它们共同决定了 filter 和 output 阶段的并发方式;当下游变慢时,背压如何从 output 一路传回 input。
核心问题:为什么把 pipeline.workers 调大,有时吞吐上升、有时毫无变化甚至更差?
pipeline 执行结构
Logstash 的 pipeline 里有两类线程,职责分开、数量独立配置:
1 | |
每个 worker 是一个独立线程,独立地从队列里拉一批 event,顺序跑完整条 filter 链,再顺序跑完整条 output 链,才算处理完这批,然后回去拉下一批。filter 链和 output 链对每个 worker 是串行的——并发来自多个 worker 同时工作,而不是单个 worker 内部并发。
上面的图给每个 worker 各画了一行 filter/output 链,但那是同一套插件实例,不是每个 worker 各持一份副本。插件的线程安全是插件自身的契约,Logstash 不替它做隔离,所以在 filter 或 output 里保存跨 event 的状态并不安全。这个前提在下面"worker 数不一定等于配置值"一节里有更直接的后果。
三个核心参数
pipeline.workers
pipeline.workers 决定 filter+output 阶段能同时跑多少个 worker 线程。默认值等于 JVM 可见的逻辑 CPU 核心数。
worker 数量不是越大越好,上界由瓶颈位置决定而不是由核数决定,具体的数字推演放在后面的背压章节。
适合调大 worker 的场景:filter 阶段本身是 CPU 密集型(大量正则匹配、Grok 解析),且 CPU 核心有空余。
不适合调大 worker 的场景:output 是 I/O 瓶颈(ES 写入延迟高、网络拥塞)或 filter 链里有阻塞操作(如 http filter 做同步外部调用)。
pipeline.batch.size
pipeline.batch.size 控制每个 worker 每次从队列里一次性拉取的最大 event 数。默认值 125(适用 Logstash 7.x / 8.x)。
batch 大小影响两件事:
- 摊薄每条 event 的调度开销:批量越大,调度开销的单条均摊越低,这对吞吐有利。
- 占用更多内存:每个 worker 同时持有一个完整 batch 的 event 对象,worker 数 × batch.size × 单条 event 内存大致是峰值内存消耗。调大 batch 前需确认 JVM 堆足够。
- 对延迟有影响:batch 越大,一批里最后一条 event 等待"凑满批"的时间越长,尾部延迟越高(见下面的 batch.delay)。
pipeline.batch.delay
pipeline.batch.delay 是凑批超时:worker 等待 batch 凑满的最长时间,默认 50 毫秒。超时到了,即使当前 batch 还没满 batch.size,worker 也会带着现有的 event 继续往下走。
低流量场景(event 进入速度比 batch.size 慢)下,这个参数决定了延迟下限。如果流量稀疏但对实时性要求高,可以把 batch.delay 调小;如果追求吞吐、能接受高延迟,可以加大。
1 | |
worker 数不一定等于配置值
pipeline.workers 是请求值,不是保证值。有两条路径会让实际生效的 worker 数与配置对不上,其中一条完全静默。
pipeline.ordered:保序把并发锁死在 1
pipeline.ordered 决定 pipeline 是否保证输出顺序与进入顺序一致,默认 auto。三种取值的语义:
| 取值 | 行为 |
|---|---|
auto(默认) |
仅当显式把 pipeline.workers 设为 1 时才启用保序 |
true |
强制保序;若 pipeline.workers > 1,启动直接失败 |
false |
关闭保序,不管 worker 数是多少 |
auto 的判定条件是"显式设置过"而不是"取值为 1"。单核机器上 pipeline.workers 的默认值恰好就是 1,但只要没在 logstash.yml 或命令行里写出来,auto 不会启用保序。想要保序,必须自己写一遍 pipeline.workers: 1。
true 那一行是硬失败而不是降级:配了 pipeline.ordered: true 又留着多 worker,Logstash 在启动阶段就报 enabling the 'pipeline.ordered' setting requires the use of a single pipeline worker 并退出。这条规则是"为什么某些 pipeline 必须 workers=1"的答案。按纯吞吐视角把这类 pipeline 的 worker 调大,等于把保序需求调坏。
非 threadsafe 的 filter 会静默把 worker 压到 1
这是"调大 workers 毫无变化"最容易漏掉的原因。Logstash 启动时按插件的 threadsafe? 标记把 filter 链分成两组,只要存在被标记为不安全的 filter(aggregate 这类跨 event 维护状态的插件即是),后续行为分两种:
- 没有显式设置
pipeline.workers:Logstash 把 worker 数强制改成 1,只在日志里留一句 warn:Defaulting pipeline worker threads to 1 because there are some filters that might not work with multiple worker threads。此时机器有多少核都不影响,实际并发就是 1。 - 显式设置了
pipeline.workers > 1:Logstash 只警告不阻止,日志里是Warning: Manual override - there are filters that might not work with multiple worker threads,并发照给,但那个不安全的 filter 可能产出错误结果。
两条路径合起来的结论是:不要假定生效的 worker 数等于 CPU 核数。实际值可以直接读出来,Node Stats 的 per-pipeline 响应里有 pipeline 这一段,写的就是运行时真正用上的三个值:
1 | |
调优前先看这一行,比按核数推算靠得住。
实验:观察 batch 行为
用最小配置复现 batch 的工作方式。准备一个 pipeline09.conf:
1 | |
用如下命令以单 worker、batch.size=3 启动:
1 | |
快速输入 5 行文本后停止输入,观察终端输出的时机:前三条会攒成一批(或等 200ms 超时)一起处理,第 4、5 条形成第二批。输出的时间间隔直接体现了"凑批 → filter 串行 → output"的节奏。
对应内部对象
上面的实验现象对应 Logstash 内部的具体机制:
| 实验现象 | 对应的内部对象 / 行为 |
|---|---|
| 5 行按 3+2 分批处理 | QueueBatch:每次从队列拉取,最多 batch.size 条 |
| 200ms 超时后未凑满也继续 | 实验跑在默认内存队列上,读端是 JrubyWrappedSynchronousQueueExt(底层 ArrayBlockingQueue);凑批等待逻辑在它与 PQ 读端的共同父类 QueueReadClientBase,字段 waitForMillis 默认 50 |
| 每批都经过同一套 filter | worker 线程由 WorkerLoopThread 包装,循环体是 org.logstash.execution.WorkerLoop;filter/output 插件实例在所有 worker 间共享 |
| 输出完才拉下一批 | worker 主循环是 readClient.readBatch() → execution.compute(...) → readClient.closeBatch(batch),三步串行,不是流水线 |
背压:从 output 到 input 的压力传递
Logstash 的 pipeline 内部是一条背压链,压力可以从尾部自然传回头部:
1 | |
这条背压链是自然的、被动的,不需要显式配置。它的存在解释了一个常见观察:ES 集群变慢,Logstash 的 input 端读取速度也跟着下降,而不是 Logstash 在本地堆积无限量的 event。
倒数第三步那个"队列满"的触发点比图上看起来靠前得多。内存队列的容量不是一个可以单独配置的参数,它由 pipeline.batch.size × pipeline.workers 算出来:8 核机器上默认是 125 × 8 = 1000 条,一千条 event 就是全部缓冲深度。max_unread_events 这个名字容易被当成内存队列的旋钮,它其实是 PQ 侧由 queue.max_events 派生出来的容量指标。所以内存队列下背压的触发几乎是即时的:下游一慢,一千条填满,input 立刻开始阻塞。要更深的缓冲只有换 PQ,靠 queue.max_bytes 和 queue.max_events 撑出量级更大的空间。
背压与 worker 数量的关系
如果 output 是瓶颈,增加 pipeline.workers 不能解决问题。考虑这个场景:ES 每批 bulk 需要 500ms,每个 worker 一批有 125 条 event,当前有 4 个 worker。每秒处理能力上限是 4 × (1000/500) × 125 = 1000 events/s。把 worker 调到 8,每秒能力变成 8 × 2 × 125 = 2000 events/s——但如果 ES 此时只能接受 1000 events/s,多出来的 4 个 worker 只是多了 4 个同时等待 ES 响应的线程,队列同样会积压,背压同样会传回 input。
真正的提速需要先找到瓶颈所在:filter CPU 密集 → 加 worker;output I/O 密集 → 调大 batch.size(减少 bulk 次数)或扩容下游。
用 Node Stats API 定位瓶颈
Logstash 提供 Node Stats API,可以直接观察各阶段耗时,不用猜测瓶颈位置:
1 | |
关注以下字段:
1 | |
字段名有两个坑要先避开。一是 per-pipeline 的 event 数就叫 queue.events,不叫 events_count;events_count 只存在于 _node/stats 的顶层聚合里,而且那个数字只累加 type == 'persisted' 的 pipeline 并跳过 system pipeline,用它看单条管道会得到 0。二是除了 queue.type 之外,上面这组 queue 字段只在 queue.type: persisted 时才注册。用默认内存队列跑本文实验时,queue.events 根本不会出现在响应里。
判读方法:如果 output 插件的 duration_in_millis / events.out 显著高于 filter 插件的,output 是瓶颈。队列侧要按同量纲比——capacity.queue_size_in_bytes 对 capacity.max_queue_size_in_bytes,queue.events 对 capacity.max_unread_events;任一组持续贴顶,说明 PQ 已经被背压填满。
直接量度背压的两个字段
上面那组队列字段回答的是"积压了多少",而背压本身是"input 被拦了多久",量它的字段是另外两个,且不依赖 PQ:
1 | |
queue_push_duration_in_millis 是累计值,单次采样的绝对数字没有意义,要看两次采样之间的增量:这个值随时间显著增长,就说明 input 线程正卡在写队列上,背压已经传到了链条头部。flow.queue_backpressure 是 Logstash 已经做完差分的速率形式,看趋势更省事。用默认内存队列时这两个字段是唯一能落地的背压观测点,因为 queue.events 那一组还没注册。
flow 命名空间下同批还有 input_throughput / filter_throughput / output_throughput / worker_concurrency,配合起来能把"哪一段慢"和"input 被拦了多久"对上。
模式提炼
1 | |
工程迁移表
| Logstash 概念 | Kafka Consumer 对应 | Flink 对应 | 通用 ETL 对应 |
|---|---|---|---|
| pipeline.workers | consumer 线程数 / partition 数 | 算子并行度 | 并发 worker 线程数 |
| batch.size | max.poll.records |
缓冲区大小 | 批量提取行数 |
| batch.delay | fetch.max.wait.ms |
缓冲超时 | 凑批超时 |
| 队列满 → input 阻塞 | producer backpressure | credit-based 背压 | 上游限速 |
| Node Stats API | Consumer group lag | Flink Web UI metrics | 作业监控指标 |
Kafka Consumer 的 max.poll.records 和 Logstash 的 batch.size 是同一个思路:每次拉取的最大消息数,用来在延迟和吞吐之间取舍。Flink 的 credit-based 背压也是同样的被动传导机制,下游处理慢则上游算子自然减速。
常见误解
误解一:“pipeline.workers 越多越好”。超过核数还有没有收益,取决于 worker 的时间花在哪里:filter 是 CPU 密集型时,核数之外的线程只是在抢同一批 CPU,收益被上下文切换吃掉;worker 大部分时间阻塞在 output 的 I/O 等待上时,超配是官方认可的做法。真正跑不掉的约束在内存侧——在途 event 的峰值是 workers × batch.size,把 worker 从 8 调到 32 就是把在途量翻两番,堆压力上来之后 GC 抖动会把并发收益整个抵掉。上界是 heap 能容纳多少在途 event,不是机器有几个核。
误解二:“batch.size 越大吞吐越高”。batch 大小增加意味着每个 worker 同时持有更多 event 对象,JVM 堆压力随之增大。当堆内存不足时,频繁的 Full GC 会导致吞吐反而下降,延迟急剧升高。调大 batch.size 必须配合 JVM 堆大小同步评估。
误解三:“背压只是 Logstash 内部的事”。背压是一条完整的链,Logstash 的 input 阻塞意味着对上游数据源的确认也会停止:Beats 的 ack 不发出、Kafka 的 offset 不提交。上游在这段时间内必须自己维持数据。这是系统级联行为,不能孤立看 Logstash。
误解四:“filter 和 output 是流水线并发的”。Logstash 的 worker 模型里,同一个 worker 的 filter 和 output 是串行的——先跑完所有 filter,再跑所有 output。不同 worker 之间是并发的,但单个 worker 内部不是流水线。如果想实现真正的 filter/output 流水线并发,需要用 Multiple Pipelines(第 12 篇)把 pipeline 拆开,用队列连接。
练习
-
用本文的
pipeline09.conf实验,分别设置pipeline.workers为 1、2、4,pipeline.batch.size为 10、50、125,观察 Node Stats API 里events.out的变化趋势。记录在哪个参数组合下吞吐开始停止增长,尝试解释原因(filter sleep 的 CPU 占用是关键提示)。 -
构造一个 output 阻塞的场景:把 output 换成
http { url => "http://localhost:9999" }(一个不存在的地址),每隔几秒采一次pipelines.main.events.queue_push_duration_in_millis,对相邻两次采样做差。验证 output 完全阻塞时,背压是否已经传回 input。顺便确认这个场景下pipelines.main.queue.events为什么取不到值。 -
在 filter 链里加一个
aggregatefilter(非 threadsafe),先不显式设置pipeline.workers启动,去_node/stats/pipelines/main读pipeline.workers的实际值,和机器核数对照。再用-w 4重跑一次,对比两次启动日志里 warn 的措辞差别,说明这两条 warn 分别意味着什么。 -
把
pipeline.ordered设成true、pipeline.workers保留 2,记录启动失败的信息。再把 workers 改成 1 确认能正常启动。最后删掉pipeline.ordered只留pipeline.workers: 1,说明默认的auto在这一步为什么同样保序,而在没写pipeline.workers的单核机器上为什么不保序。
系列导航
参考资料
- Logstash pipeline 配置文档:https://www.elastic.co/guide/en/logstash/current/logstash-settings-file.html(pipeline.workers、batch.size、batch.delay 参数说明)
- Logstash Node Stats API:https://www.elastic.co/guide/en/logstash/current/node-stats-api.html(pipeline、event、plugin、queue 各维度指标)
- Logstash 性能调优文档:https://www.elastic.co/guide/en/logstash/current/performance-tuning.html(worker 和 batch 调优的官方指导)
- Logstash 源码仓库:https://github.com/elastic/logstash(
logstash-core/lib/logstash/java_pipeline.rb里的safe_pipeline_worker_count与preserve_event_order?、logstash-core/src/main/java/org/logstash/execution/WorkerLoop.java、同目录QueueBatch.java) - Kafka Consumer 配置文档:https://kafka.apache.org/documentation/#consumerconfigs(max.poll.records、fetch.max.wait.ms,与 batch.size / batch.delay 对照)
- Flink 背压机制文档:https://nightlies.apache.org/flink/flink-docs-stable/docs/ops/monitoring/back_pressure/(credit-based 背压与 Logstash 被动背压的对比)
