深入 Logstash 13 - 监控:Node Stats API、hot threads 与瓶颈定位
上一篇解决了 Multiple Pipelines 的隔离与路由问题。这一篇进入运维调优阶段的第一个主题:怎么从指标判断管道卡在 input、filter 还是 output,以及 hot threads API 的读法。
核心问题:一条管道吞吐下降时,靠什么指标定位瓶颈在哪一段?plugin 级别的耗时从哪里拿到?hot threads 输出能读出什么信息?
监控数据的来源
1 | |
Logstash 自带 HTTP API,进程启动后立即可用,不需要额外安装任何组件,也不依赖 X-Pack 授权。官方目前列出五个监控 API:
1 | |
最后一个是专做健康与瓶颈判定的,本篇后面的判断树是手工版本,_health_report 是官方版本,两者可以互相印证。
接监控时第一批会撞上的是端口和绑定地址这两件事,它们的默认值都和"硬编码 localhost:9600"的直觉不一致:
1 | |
后两行是改 api.http.host 之前要一起看的:这个 API 默认零鉴权零 TLS。"进程启动后立即可用"的另一半,是一个默认对外无鉴权的运维端口。对外暴露前要配上:
1 | |
密码类字段用 keystore 或环境变量占位,避免明文落在配置文件里。
指标送出去有三条路线,本篇后面只展开第一条:
1 | |
只需要临时看一眼或者写脚本告警,用第一条。要把指标接进既有的可观测性栈,第二条是方向正确的那条路,但它还在预览阶段,生产上要先确认版本够新、并接受接口可能变动。
关键对象与数据结构
Node Stats API 的响应按几个顶层 section 组织:
1 | |
queue 这一段的层级最容易记错,而后面的判断树全靠它:
1 | |
比字段名更要紧的是一条前提:整个 queue 段只在 queue.type: persisted 时才注册。默认配置是内存队列,照上面的路径去读会直接拿不到字段,这既不是版本问题也不是权限问题。
内存队列的积压得换两个间接信号来看:
1 | |
events.in、events.filtered、events.out 和各插件的 duration_in_millis 都是自启动以来的累计值,不是实时速率。算速率要么两次采样做差除时间间隔,要么直接读 flow——后者不只是省掉手工差分,它多出来的 worker_utilization 是现在判读瓶颈的主线,下一节就从它开始。flow metrics 从 8.5.0 起进入 Node Stats API,PQ 的两个增长率指标和 worker_utilization 更晚。
瓶颈定位逻辑
先说清一件决定后面所有读法的事:pipeline.workers 控制的那组线程同时执行 filter 和 output 两个阶段,不存在独立的 output 线程池。WorkerLoop 的类注释写得很直白——它负责为每个 batch 执行 filter 和 output 插件。
这一条直接解释了几个否则很难串起来的现象:output 阻塞时 filter 也一起停(同一个线程卡在 output 里,回不到 filter);给 pipeline.workers 加值同时抬高 filter 并发和 output 并发;以及第 12 篇里 output isolator 拓扑为什么必要——想让两个 output 互不牵连,只能把它们放进不同 pipeline,因为同一条 pipeline 里它们连线程都是共用的。
首选路径:flow metrics
1 | |
插件级的 worker_utilization 是这条路径的关键:pipeline 级说明 worker 满了,插件级说明满在谁身上。因为 filter 和 output 共用 worker,这张表不需要事先区分是 filter 慢还是 output 慢——占比最大的那个插件是谁就是谁。
兜底路径:手工差分累计值
flow metrics 之前的版本,或者需要更细的自定义口径时,回到累计值做差:
1 | |
单个插件的耗时从 plugins.filters[i].events.duration_in_millis 读取。把同一插件两次采样的 duration_in_millis 差值除以对应事件数差值,得到该插件的平均每事件处理时间(毫秒/event)——这是插件级 worker_millis_per_event 的手工版本。
实验:用 curl 读 Node Stats
启动一个本地 Logstash 实例,执行以下命令:
1 | |
两次采样、做差,得到速率:
1 | |
如果版本支持 flow metrics,上面这套差分都不用写,一次请求就能拿到速率和饱和度:
1 | |
worker_utilization 接近 100 就说明 worker 已经饱和,这是下一节判断树的入口。
hot threads API
_node/hot_threads 返回当前 CPU 占用最高的 Java 线程的堆栈快照。默认返回 JSON,加上 human=true 才是纯文本——human 这个通用参数在 Logstash 里只对 hot threads API 生效:
1 | |
输出模板写死在 logstash-core/locales/en.yml 里,只有两行结构:主机名标题行 + Hot threads at ... 一行,然后每个线程一行摘要加它的栈:
1 | |
每行摘要有四项:CPU 百分比、线程状态、线程名、线程 id。没有 interval= 字段,也没有 “(Xms out of Yms)” 这种窗口占比——原因见下面对百分比语义的说明。
线程名不是随手起的,Linux 平台上 Logstash 给线程打的标签有固定规律:
1 | |
多条 pipeline 时,方括号里就是各自的 pipeline.id,这是把热点线程归属到具体管道的唯一线索。
合法查询参数只有四个(外加通用的 human):
1 | |
没有 interval 参数。写 ?interval=1000 不会报错,也不会有任何效果——参数被直接忽略,读者以为自己调过了采样窗口,实际什么都没发生。
之所以没有采样窗口,是因为百分比根本不是窗口采样算出来的。它的算法是线程累计 cpu.time 除以进程 uptime,也就是"这个线程自 Logstash 启动以来占掉了多少 CPU"。
这个区别会直接改变结论的可信度。累计值意味着一个只在启动阶段疯狂跑过的线程,会长期挂在榜首;而一个刚刚开始变慢的线程,要很久才能爬上来。所以"多次调用、看同一个线程是否持续出现"这种做法在这里是无效的——它区分不了"现在正热"和"开机时热过",因为两者在累计视角下都表现为持续出现。
要判断当下的热点,得自己取两次快照做差。JSON 响应里给的是 percent_of_cpu_time,把它乘回同一时刻的 jvm.uptime_in_millis 就还原成累计 CPU 毫秒数,两次相减即为窗口内的真实消耗:
1 | |
30 秒窗口内消耗 CPU 最多的线程才是当前的热点。这一步是 hot threads 从"看个大概"变成可用证据的分界线。
将实验结果映射回内部对象
Node Stats 的数字直接对应管道内部的对象。这些类名都能在 elastic/logstash 仓库里打开对应文件:
1 | |
hot threads 里的 [<id>]>workerN 线程就是 WorkerLoopThread,循环体是 WorkerLoop,数量与该 pipeline 的 pipeline.workers 一一对应。它在一次循环里读一个 batch、跑完 filter、再交给 output delegator,所以 filter 和 output 的耗时都记在同一个线程的账上。
output 侧的并发因此不是"再开一组线程",而是插件自己的连接池。以 ES output 为例,真正的旋钮是 pool_max(默认 1000)和 pool_max_per_route(默认 100)。
这里有一个值得单独记住的陷阱:output plugin 有一个 workers 参数,写进配置不会报错,但它在 logstash-core/lib/logstash/outputs/base.rb 里的声明是 :deprecated => "This parameter will be ignored."。参数存在、管道正常启动、只在日志里留一条 deprecation,然后被完全忽略。这比一个不存在的参数更难发现——不存在的参数会让管道起不来,读者立刻知道;这个会让读者以为自己已经调过了 output 并发,实际什么都没发生。
模式提炼
1 | |
这套由粗到细的路径不依赖 Logstash 特有的实现。当全局饱和度、分段耗时、线程栈这三类数据齐备时,都可以这样排查。
工程迁移表
| Logstash 监控概念 | Kafka 生态对应 | Flink 对应 | 通用 ETL 对应 |
|---|---|---|---|
flow.*_throughput 预计算速率 |
consumer/producer rate(JMX) | records-in-rate / records-out-rate | 行/记录吞吐量 |
flow.worker_utilization |
consumer 线程忙闲比 | Task busyTimeMsPerSecond | Worker 池饱和度 |
per-plugin duration_in_millis |
UDF 执行时间(JMX metrics) | operator latency histogram | Transform 阶段耗时 |
queue.events 增长(PQ) |
consumer group lag 增长 | checkpoint 间隔内积压量 | staging 表行数增长 |
queue_push_duration_in_millis(内存队列) |
producer send 阻塞时间 | 反压时的 buffer 等待 | 写入阻塞耗时 |
hot threads [main]>workerN 栈位置 |
consumer poll / producer send 线程栈 | TaskManager thread dump | Worker 线程栈 |
jvm.heap_used_percent + GC 曲线形态 |
broker/consumer heap OOM 风险 | TaskManager GC 压力 | JVM 服务 GC 压力 |
常见误解
误解一:“events.out 低就是 output 的问题”。它只说明管道整体输出慢。因为 filter 和 output 跑在同一组 worker 线程上,这两种情况在 events.out 上表现完全一样:filter 自己慢,或者 output 卡住把整个 worker 一起拖住。区分它们要看插件级的 worker_utilization(或两次采样的 per-plugin duration_in_millis),再加上 input 本身是否就慢(源数据量小)。
误解二:“hot threads 里 worker 线程 CPU 高就该加 worker”。CPU 高可能是单个 worker 执行效率低(比如灾难性回溯的 Grok 正则),加 worker 只是把同样的浪费摊到更多线程上,不解决根因。而且 hot threads 的百分比是自启动以来的累计占比,看到某个线程排在前面并不等于它现在正忙。先对两次快照做差确认它当下确实在烧 CPU,再用 Grok Debugger 确认正则是否有指数级回溯,最后才轮到加 worker。
误解三:“Node Stats 的 queue 字段拿不到值是版本或权限问题”。queue 段只在 queue.type: persisted 时注册。默认是内存队列,此时这条 pipeline 下面根本没有 queue 对象,读到的是空。内存队列要看 events.queue_push_duration_in_millis 或 flow.queue_backpressure。同理,节点顶层的 queue.events_count 是个只统计 PQ pipeline 的 rollup,全用内存队列时它恒为 0。
练习
-
本地启动一条带慢 filter 的 Logstash 管道(用
sleep { time => 0.1 }模拟),向它发送事件,每隔 5 秒采集一次_node/stats,计算events.out的速率,确认与 sleep 设置的延迟吻合(约 10 events/s)。再把pipeline.workers从 1 调到 2,重复实验,观察速率变化。 -
对同一个带慢 filter 的管道调用
_node/hot_threads?human=true,确认线程名形如[main]>worker0(多条 pipeline 时方括号里是各自的pipeline.id),并确认堆栈指向 sleep 调用。把 sleep 换成 CPU 密集型操作(比如循环 Grok),重复观察堆栈结构的变化。顺便试一次?interval=1000,对比加与不加的输出,验证这个参数被静默忽略。 -
用上面那段两次快照做差的脚本,在管道空转 5 分钟后再制造一次短促的 CPU 尖峰。对比"累计百分比排名"和"30 秒窗口内 CPU 增量排名"两个榜单,找出两者给出不同答案的具体线程,解释为什么单看累计值会把结论带偏。
-
构造一个 output 慢的场景(用
httpoutput 指向一个延迟高的端点)。先用默认内存队列跑一遍,确认pipelines.<id>.queue整段不存在,只能靠events.queue_push_duration_in_millis判断积压;再把queue.type改成persisted重跑,此时queue.events与queue.capacity.queue_size_in_bytes才有值。结合插件级worker_utilization确认是 output 而非 filter 造成的积压。
系列导航
参考资料
- Logstash Monitoring API 文档:https://www.elastic.co/guide/en/logstash/current/monitoring-logstash.html(五个监控 API 清单、
api.ssl.*与api.auth.*加固、human参数只对 hot threads 生效) - Logstash Node Stats API 参考:https://www.elastic.co/guide/en/logstash/current/node-stats-api.html(events、pipeline、jvm、process 各 section)
- Logstash Flow Metrics:https://www.elastic.co/guide/en/logstash/current/flow-metrics.html(
input_throughput/filter_throughput/output_throughput/queue_backpressure/worker_concurrency/worker_utilization,以及 PQ 的queue_persisted_growth_events/queue_persisted_growth_bytes) - Logstash 性能排查与调优:https://www.elastic.co/guide/en/logstash/current/tuning-logstash.html(以
worker_utilization为主线的判读顺序,以及 Linux 下的线程命名规则[base]>workerN/[base]<inputname) - Logstash Pipeline Viewer(Kibana):https://www.elastic.co/guide/en/logstash/current/logstash-pipeline-viewer.html(可视化拓扑与 plugin 级别指标)
- Logstash 源码(hot threads 输出模板):https://github.com/elastic/logstash/blob/main/logstash-core/locales/en.yml(
logstash.web_api.hot_threads两行模板,标题行与每线程摘要行的确切字段) - Logstash 源码(hot threads 百分比算法):https://github.com/elastic/logstash/blob/main/logstash-core/lib/logstash/api/commands/hot_threads_reporter.rb(
cpu_time_as_percent用cpu.time / uptime,是累计占比) - Logstash 源码(pipeline 指标注册):https://github.com/elastic/logstash/blob/main/logstash-core/src/main/java/org/logstash/execution/AbstractPipelineExt.java(
queue.capacity.*与 PQ 专属指标的注册条件) - Logstash 源码(output workers 已废弃):https://github.com/elastic/logstash/blob/main/logstash-core/lib/logstash/outputs/base.rb(
config :workers, :deprecated => "This parameter will be ignored.")
