上一篇(08)覆盖了 Elasticsearch output 的 bulk 写入机制。这一篇进入 pipeline 执行模型本身:pipeline.workerspipeline.batch.sizepipeline.batch.delay 三个参数各自控制什么,它们共同决定了 filter 和 output 阶段的并发方式;当下游变慢时,背压如何从 output 一路传回 input。

核心问题:为什么把 pipeline.workers 调大,有时吞吐上升、有时毫无变化甚至更差?

pipeline 执行结构

Logstash 的 pipeline 里有两类线程,职责分开、数量独立配置:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
input 线程(1 条或多条,取决于 input 插件)
│ 每条 event 写入队列

┌──────────────────────────────┐
│ queue(内存 / PQ) │
└──────────────────────────────┘
│ 每次拉取一批(batch)

worker 0 ──▶ filter chain ──▶ output chain
worker 1 ──▶ filter chain ──▶ output chain
worker 2 ──▶ filter chain ──▶ output chain
...
worker N ──▶ filter chain ──▶ output chain
(N = pipeline.workers,默认等于 CPU 逻辑核心数)

每个 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
2
3
4
# logstash.yml 相关参数示意
pipeline.workers: 4
pipeline.batch.size: 125
pipeline.batch.delay: 50

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
2
3
curl -s http://localhost:9600/_node/stats/pipelines/main \
| python3 -c 'import json,sys; print(json.load(sys.stdin)["pipelines"]["main"]["pipeline"])'
# {'workers': 1, 'batch_size': 125, 'batch_delay': 50}

调优前先看这一行,比按核数推算靠得住。

实验:观察 batch 行为

用最小配置复现 batch 的工作方式。准备一个 pipeline09.conf

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
input {
stdin {
codec => line
}
}

filter {
sleep {
time => "0.1" # 模拟 100ms filter 处理耗时
every => 1
}
}

output {
stdout {
codec => rubydebug
}
}

用如下命令以单 worker、batch.size=3 启动:

1
2
3
4
bin/logstash -f pipeline09.conf \
--pipeline.workers 1 \
--pipeline.batch.size 3 \
--pipeline.batch.delay 200

快速输入 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
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
ES 写入变慢


output.elasticsearch 的 HTTP 请求阻塞(等待 ES 响应)


worker 线程卡在 output 阶段,无法返回到拉取下一批的状态


队列里已有的 event 积压(没有 worker 消费)


队列满(内存队列容量 = batch.size × workers;PQ 受 queue.max_bytes 与 queue.max_events 限制)


input 线程调用 queue.write() 时阻塞


input 对上游的读取操作减速(如 Beats input 不发 ACK,Kafka input 不提交 offset)

这条背压链是自然的、被动的,不需要显式配置。它的存在解释了一个常见观察: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_bytesqueue.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
2
# 查看 pipeline 级别的统计
curl -s http://localhost:9600/_node/stats/pipelines/main | python3 -m json.tool

关注以下字段:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
pipelines.main.events.in              # 进入 pipeline 的总 event 数
pipelines.main.events.out # 离开 pipeline 的总 event 数
pipelines.main.events.filtered # 经过 filter 的数
pipelines.main.events.duration_in_millis # 所有 event 在 pipeline 的总耗时

# 每个 plugin 的耗时(filter/output 各自)
pipelines.main.plugins.filters[*].events.duration_in_millis
pipelines.main.plugins.outputs[*].events.duration_in_millis

# 队列状态
pipelines.main.queue.type # memory 或 persisted
pipelines.main.queue.events # 当前队列里未读的 event 数
pipelines.main.queue.capacity.queue_size_in_bytes # 队列占用的字节数
pipelines.main.queue.capacity.max_queue_size_in_bytes # 即 queue.max_bytes
pipelines.main.queue.capacity.max_unread_events # 由 queue.max_events 派生

字段名有两个坑要先避开。一是 per-pipeline 的 event 数就叫 queue.events,不叫 events_countevents_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_bytescapacity.max_queue_size_in_bytesqueue.eventscapacity.max_unread_events;任一组持续贴顶,说明 PQ 已经被背压填满。

直接量度背压的两个字段

上面那组队列字段回答的是"积压了多少",而背压本身是"input 被拦了多久",量它的字段是另外两个,且不依赖 PQ:

1
2
pipelines.main.events.queue_push_duration_in_millis  # input 在 queue.write() 上累计阻塞的毫秒数
pipelines.main.flow.queue_backpressure # 同一数据的速率形式(flow 指标)

queue_push_duration_in_millis 是累计值,单次采样的绝对数字没有意义,要看两次采样之间的增量:这个值随时间显著增长,就说明 input 线程正卡在写队列上,背压已经传到了链条头部。flow.queue_backpressure 是 Logstash 已经做完差分的速率形式,看趋势更省事。用默认内存队列时这两个字段是唯一能落地的背压观测点,因为 queue.events 那一组还没注册。

flow 命名空间下同批还有 input_throughput / filter_throughput / output_throughput / worker_concurrency,配合起来能把"哪一段慢"和"input 被拦了多久"对上。

模式提炼

1
2
3
4
5
6
7
模式:批量拉取 + worker 并发 + 被动背压

- 把 filter+output 拆进独立 worker,用 worker 数量控制并发度
- 每个 worker 一次拉取一批,用 batch.size 摊薄调度开销
- 用 batch.delay 控制延迟与吞吐的取舍点
- 不为背压单独设计,让队列容量成为天然的流量控制机制
- 背压自动从 output 传回 input,防止无限制的内存积累

工程迁移表

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 拆开,用队列连接。

练习

  1. 用本文的 pipeline09.conf 实验,分别设置 pipeline.workers 为 1、2、4,pipeline.batch.size 为 10、50、125,观察 Node Stats API 里 events.out 的变化趋势。记录在哪个参数组合下吞吐开始停止增长,尝试解释原因(filter sleep 的 CPU 占用是关键提示)。

  2. 构造一个 output 阻塞的场景:把 output 换成 http { url => "http://localhost:9999" }(一个不存在的地址),每隔几秒采一次 pipelines.main.events.queue_push_duration_in_millis,对相邻两次采样做差。验证 output 完全阻塞时,背压是否已经传回 input。顺便确认这个场景下 pipelines.main.queue.events 为什么取不到值。

  3. 在 filter 链里加一个 aggregate filter(非 threadsafe),先不显式设置 pipeline.workers 启动,去 _node/stats/pipelines/mainpipeline.workers 的实际值,和机器核数对照。再用 -w 4 重跑一次,对比两次启动日志里 warn 的措辞差别,说明这两条 warn 分别意味着什么。

  4. pipeline.ordered 设成 truepipeline.workers 保留 2,记录启动失败的信息。再把 workers 改成 1 确认能正常启动。最后删掉 pipeline.ordered 只留 pipeline.workers: 1,说明默认的 auto 在这一步为什么同样保序,而在没写 pipeline.workers 的单核机器上为什么不保序。

系列导航

序号 主题
00 导读:核心对象是 event,骨架是三段管道
01 架构:JRuby、JVM 与 pipeline 的运行形态
02 event 模型:@timestamp、@metadata 与字段引用
03 codec:字节流与 event 的边界转换
04 input 插件:拉取、监听与 Beats 接入
05 Grok 的本质:命名正则加预定义 pattern
06 dissect 与结构化 filter:放弃回溯换吞吐
07 常用 filter 组合:mutate、date、geoip 与条件
08 output 插件:Elasticsearch output 与批量写入
09 pipeline 执行模型:worker、batch 与背压(本篇)
10 内存队列 vs 持久队列:可靠性的分界线
11 死信队列(DLQ):无法处理的 event 去哪
12 Multiple Pipelines 与 pipeline-to-pipeline
13 监控:Node Stats API、hot threads 与瓶颈定位
14 性能调优:JVM heap、批处理与持久队列磁盘
15 Logstash vs Beats vs Ingest Pipeline:该用谁
16 Logstash vs Fluentd vs Vector:日志管道的三种取舍
17 Logstash 的演进与 Elastic Agent 的冲击

参考资料