上一篇(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 内部并发。

三个核心参数

pipeline.workers

pipeline.workers 决定 filter+output 阶段能同时跑多少个 worker 线程。默认值等于 JVM 可见的逻辑 CPU 核心数。

worker 数量不是越大越好。每个 worker 都是一个线程,线程本身有调度和内存开销;当 output 下游是瓶颈时(比如 ES 写入已经饱和),增加 worker 只会让更多线程同时卡在 output 的等待上,而不会提升整体吞吐。

适合调大 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

实验:观察 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 超时后未凑满也继续 WrappedAckedQueue#readBatch 的 poll timeout 逻辑
每批都经过同一套 filter 同一个 PipelineThread(worker)持有 filter chain 引用,顺序执行
输出完才拉下一批 worker 的主循环:process_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 消费)


队列满(内存队列 maxUnread 限制 / PQ max_bytes 限制)


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


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

这条背压链是自然的、被动的,不需要显式配置。它的存在解释了一个常见观察:ES 集群变慢,Logstash 的 input 端读取速度也跟着下降,而不是 Logstash 在本地堆积无限量的 event。

背压与 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
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.events_count # 当前队列里的 event 数
pipelines.main.queue.queue_size_in_bytes

如果 output 插件的 duration_in_millis / events.out 显著高于 filter 插件的,output 是瓶颈。如果 queue 的 events_count 持续接近上限,队列已经在被背压充满。

模式提炼

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 增加到超过 CPU 核心数后,线程切换的开销会抵消并发收益;更重要的是,如果瓶颈在 output 侧(ES 写入饱和),增加 worker 只会让更多线程同时阻塞等待,不提升吞吐,反而会因内存占用增加而加剧 GC 压力。

误解二:“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" }(一个不存在的地址),观察 Node Stats 里 queue.events_count 的变化趋势。验证当 output 完全阻塞时,背压是否从队列传回了 input。

  3. 查阅 Logstash 8.x 和 7.x 的官方文档,确认 pipeline.batch.size 的默认值是否一致。如果存在版本差异,说明差异对"用默认值运行 Logstash"的行为有什么影响。

系列导航

序号 主题 状态
00 导读:核心对象是 event,骨架是三段管道 已发布
01 Logstash 架构:JRuby、JVM 与 pipeline 的运行形态
02 event 模型:@timestamp@metadata 与字段引用
03 codec:字节流与 event 的边界转换
04 input 接入:file、beats、kafka 的读取边界
05 Grok:命名正则与预置 pattern
06 dissect:分隔符场景下的 Grok 替代
07 常用 filter 组合:mutate、date、geoip
08 Elasticsearch output:bulk 写入与索引路由
09 pipeline 执行模型:worker、batch 与背压 本篇
10 内存队列 vs 持久队列:可靠性的分界线 下一篇
11 死信队列(DLQ):处理失败 event 的最后一道闸
12 Multiple Pipelines:隔离、解耦与扇出
13 监控与调优:Node Stats 与瓶颈定位
14 性能调优:JVM、batch、codec 与 filter 优化
15-17 演进、生态与对比

参考资料