上一篇解决了 Multiple Pipelines 的隔离与路由问题。这一篇进入运维调优阶段的第一个主题:怎么从指标判断管道卡在 input、filter 还是 output,以及 hot threads API 的读法。

核心问题:一条管道吞吐下降时,靠什么指标定位瓶颈在哪一段?plugin 级别的耗时从哪里拿到?hot threads 输出能读出什么信息?

监控数据的来源

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
      Logstash 进程
┌──────────────────────────────────────────────┐
│ │
│ input plugins │
│ └─ events.in counter │
│ │
│ [queue] │
│ └─ queue.events / queue_size_in_bytes │
│ │
│ filter workers ───────────────────────── │
│ └─ per-plugin duration_in_millis │
│ │
│ output plugins │
│ └─ events.out / duration_in_millis │
│ │
│ JVM / OS metrics │
│ └─ heap_used / gc / open_file_descriptors │
└──────────────────────────────────────────────┘

▼ HTTP API (默认端口 9600)
GET /_node/stats
GET /_node/hot_threads

Logstash 自带 HTTP API,默认监听 9600 端口。这个 API 不需要额外安装任何组件,进程启动后立即可用。全部监控数据通过它暴露,不依赖 X-Pack 授权。

关键对象与数据结构

Node Stats API 的响应按几个顶层 section 组织:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
GET http://localhost:9600/_node/stats

{
"events": {
"in": <total events ingested>,
"filtered": <total events passed through filters>,
"out": <total events successfully sent to outputs>,
"duration_in_millis": <cumulative filter+output processing time>,
"queue_push_duration_in_millis": <time spent pushing into queue>
},
"pipelines": {
"<pipeline_id>": {
"events": { ... }, // per-pipeline counters
"plugins": {
"inputs": [ { "id": "...", "events": { "out": N }, ... } ],
"filters": [ { "id": "...", "events": { "in": N, "out": N,
"duration_in_millis": N } } ],
"outputs": [ { "id": "...", "events": { "in": N, "out": N,
"duration_in_millis": N } } ]
},
"queue": {
"events_count": <events currently in queue>,
"queue_size_in_bytes": <current disk/memory usage>,
"max_queue_size_in_bytes": <configured cap>
}
}
},
"jvm": {
"heap_used_in_bytes": N,
"heap_used_percent": N,
"gc": { "collectors": { "old": { "collection_time_in_millis": N } } }
},
"process": {
"open_file_descriptors": N,
"cpu": { "percent": N }
}
}

events.inevents.filteredevents.out 是自启动以来的累计值,不是实时速率。计算速率需要两次采样做差除时间间隔。Logstash 8.x 引入了 flow metrics,在 pipelines.<id>.flow 下直接暴露速率和百分比,省去手动做差的步骤。

瓶颈定位逻辑

把三段的 duration_in_millisevents.out 对照来看,能直接定位慢在哪里。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
瓶颈判断树
─────────────────────────────────────────
events.in >> events.out
且 queue.events_count 持续增长
→ filter 或 output 是瓶颈

events.out 稳定
且 filter 各插件 duration 总和低
且 output duration_in_millis 高
→ output 是瓶颈(下游慢,如 ES bulk 延迟)

events.out 稳定
且某个 filter plugin duration_in_millis 独占大头
→ 该 filter 是瓶颈(通常是复杂 Grok 或 sleep)

events.in 低
且 queue 未满
→ input 是瓶颈(数据源慢,或 input plugin 限速)
─────────────────────────────────────────

单个插件的耗时从 plugins.filters[i].events.duration_in_millis 读取。把同一插件两次采样的 duration_in_millis 差值除以对应事件数差值,得到该插件的平均每事件处理时间(毫秒/event)。这是定位"哪个 filter 最慢"的标准操作。

实验:用 curl 读 Node Stats

启动一个本地 Logstash 实例,执行以下命令:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
# 查看整体 events 计数
curl -s http://localhost:9600/_node/stats | python3 -m json.tool | grep -A6 '"events"'

# 查看单条 pipeline 的 plugin 级别 duration
curl -s 'http://localhost:9600/_node/stats?pretty' \
| python3 -c "
import json, sys
data = json.load(sys.stdin)
for pid, pdata in data['pipelines'].items():
print(f'pipeline: {pid}')
for f in pdata['plugins']['filters']:
ms = f['events'].get('duration_in_millis', 0)
cnt = f['events'].get('in', 1)
print(f' filter {f[\"id\"]}: {ms}ms total, {ms/max(cnt,1):.2f}ms/event')
for o in pdata['plugins']['outputs']:
ms = o['events'].get('duration_in_millis', 0)
cnt = o['events'].get('in', 1)
print(f' output {o[\"id\"]}: {ms}ms total, {ms/max(cnt,1):.2f}ms/event')
"

两次采样、做差,得到速率:

1
2
3
4
5
# 采样间隔 10 秒,计算事件吞吐速率
T1=$(curl -s http://localhost:9600/_node/stats | python3 -c "import json,sys; d=json.load(sys.stdin); print(d['events']['out'])")
sleep 10
T2=$(curl -s http://localhost:9600/_node/stats | python3 -c "import json,sys; d=json.load(sys.stdin); print(d['events']['out'])")
echo "events/sec = $(( (T2 - T1) / 10 ))"

如果 Logstash 版本支持 flow metrics(8.x),直接读 flow.output_throughput.current 字段,已经是每秒速率,无需手动差分。

hot threads API

_node/hot_threads 返回当前 CPU 占用最高的 Java 线程的堆栈快照:

1
curl -s 'http://localhost:9600/_node/hot_threads?human=true'

典型输出结构:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
::: {logstash} ...
Hot threads at 2026-08-06T07:10:00.000Z, interval=500ms, busiestThreads=3

55.0% (275.0ms out of 500ms) cpu usage by thread
'LogStash::FilterWorker #0'
...
java.lang.Thread.sleep(Native Method)
org.jruby.ext.thread.SleepTask2.run(...)
...

30.0% (150.0ms out of 500ms) cpu usage by thread
'LogStash::OutputWorker'
...
org.elasticsearch.client.RestClient.performRequest(...)
...

读法要点:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
线程名模式          含义
───────────────────────────────────────────────────────
FilterWorker #N filter 阶段 worker 线程;如果全都高
CPU 且堆栈卡在某个 filter plugin → 该 filter 是 CPU 瓶颈
如果堆栈卡在 I/O wait / sleep → 等待下游或被限速

OutputWorker output 线程;堆栈深入网络调用且 CPU 高
→ output plugin 在等待下游响应(如 ES bulk)

pipeline-input-worker input 线程;卡在网络 read 是正常的;
如果 CPU 高且堆栈在 JSON 解析 → input codec 是瓶颈

[ruby-1]-pipeline JRuby 解释器线程;出现在 hot threads
通常意味着 Ruby 层逻辑执行密集(如复杂条件判断)

hot threads 的采样区间默认 500ms,可通过 ?threads=10&interval=1000 参数调整。多次调用、观察同一个线程是否持续出现,比单次快照更可靠。

将实验结果映射回内部对象

Node Stats 的数字直接对应管道内部的计数器对象:

1
2
3
4
5
6
7
8
API 路径                                     内部对象
───────────────────────────────────────────────────────────────────
pipelines.<id>.events.in InputQueue 的入队计数
pipelines.<id>.events.out OutputDelegator 的确认计数
pipelines.<id>.plugins.filters[i].duration FilterWorker 对该插件的累计调用时间
pipelines.<id>.queue.events_count Queue 的当前深度(PQ 或内存队列)
jvm.heap_used_percent JVM 堆占用率,超过 85% 触发频繁 GC
process.cpu.percent 进程级 CPU,包含所有线程

hot threads 里的 FilterWorkerpipeline.workers 配置的数量一一对应;OutputWorker 通常每个 output plugin 实例一个线程。filter 段的并发度由 pipeline.workers 控制,output 端的并发由各 output plugin 自己的连接池和 workers 参数控制(两者互相独立)。

模式提炼

1
2
3
4
5
6
7
8
9
10
11
12
模式:三段分层指标 + 差分速率 + 栈快照三角定位

- 第一层(全局速率):用 events.in 与 events.out 的差分
判断管道整体是否有积压

- 第二层(分段耗时):用 per-plugin duration_in_millis / events.in
定位慢在哪个 filter 或哪个 output

- 第三层(线程状态):用 hot_threads 判断是 CPU 密集还是 I/O 等待
CPU 密集 → 算法层面优化 filter;I/O 等待 → 调整 output 并发/批量

三层合用,排除方向是由粗到细:先看全局,再看分段,再看线程。

这个三角定位模式不依赖 Logstash 特有的实现。任何有计数器、耗时统计和线程 dump 的流处理系统(Flink、Kafka Streams)都遵循相同的定位路径。

工程迁移表

Logstash 监控概念 Kafka 生态对应 Flink 对应 通用 ETL 对应
events.in / events.out 差分速率 consumer lag(offset 差值) records-in-rate / records-out-rate 行/记录吞吐量
per-plugin duration_in_millis UDF 执行时间(JMX metrics) operator latency histogram Transform 阶段耗时
queue.events_count 增长 consumer group lag 增长 checkpoint 间隔内积压量 staging 表行数增长
hot threads FilterWorker CPU 高 consumer poll 线程 CPU TaskManager thread CPU Worker 线程 CPU
hot threads OutputWorker I/O wait producer send 阻塞 Sink Connector I/O 等待 Load 阶段网络等待
jvm.heap_used_percent > 85% broker/consumer heap OOM 风险 TaskManager GC 压力 JVM 服务 GC 压力

常见误解

误解一:“events.out 低就是 output 的问题”。events.out 低只说明管道整体输出慢,原因可能是 filter 太慢(事件还没处理完)、队列积压(filter 快但 output 慢导致 filter worker 也被阻塞)或 input 本来就慢(源数据量小)。需要结合 queue.events_count 和 per-plugin duration 才能区分。

误解二:“hot threads 里 FilterWorker CPU 高就一定要加 worker”。CPU 高可能是单个 worker 执行效率低(比如灾难性回溯的 Grok 正则),加 worker 只是分摊,不解决根因。先用 Grok Debugger 确认正则是否有指数级回溯,再决定是改正则还是加 worker。

误解三:“Node Stats 的数字是瞬时值”。所有 events 计数都是自启动以来的单调递增累计值。直接比较两次采样的绝对值没有意义,必须做差分除时间间隔才能得到速率。flow metrics(8.x)是例外,它已经是窗口内的速率。

误解四:“看 Kibana 的 pipeline viewer 比 curl API 更准确”。Kibana 的 pipeline viewer 底层调用的是同一套 Node Stats API,数据来源相同。API 的优势在于可以写脚本自动化采集和告警,Kibana 的优势在于可视化拓扑和颜色高亮;两者互补,不存在准确性差异。

练习

  1. 本地启动一条带慢 filter 的 Logstash 管道(用 sleep { time => 0.1 } 模拟),向它发送事件,每隔 5 秒采集一次 _node/stats,计算 events.out 的速率,确认与 sleep 设置的延迟吻合(约 10 events/s)。再把 pipeline.workers 从 1 调到 2,重复实验,观察速率变化。

  2. 对同一个带慢 filter 的管道调用 _node/hot_threads,观察 FilterWorker 线程是否出现在 hot threads 列表中,以及堆栈是否指向 sleep 调用。把 sleep 换成 CPU 密集型操作(比如循环 Grok),重复观察堆栈结构的变化。

  3. 构造一个 output 慢的场景(用 http output 指向一个延迟高的端点),观察 queue.events_count 是否持续增长,以及 hot threads 里 output 相关线程的堆栈位置。结合 per-plugin duration 确认是 output 而非 filter 造成的积压。

系列导航

序号 主题 状态
00 导读:核心对象是 event,骨架是三段管道 已发布
01 Logstash 架构:JRuby、JVM 与 pipeline 的运行形态 已发布
02-08 核心抽象与插件三段 已发布
09-12 管道执行与可靠性 已发布
13 监控:Node Stats API、hot threads 与瓶颈定位 本篇
14 性能调优:JVM heap、批处理与持久队列磁盘 下一篇
15-17 演进、生态与对比 后续阶段

参考资料