深入 Logstash 14 - 性能调优:JVM heap、批处理与持久队列磁盘
上一篇解决了怎么用 Node Stats API 和 hot threads 定位瓶颈在哪一段。这一篇进入调优执行层:拿到瓶颈定位结论之后,heap 大小、GC 策略、batch 与 worker 组合、持久队列磁盘 I/O 各应该怎么调。
核心问题:heap 大小与 GC 停顿之间的取舍;pipeline.workers 与 pipeline.batch.size 的组合效果;持久队列磁盘成为瓶颈时的判断和处置。
调优对象全景
1 | |
调优的正确顺序:先用第 13 篇的指标定位瓶颈所在层,再在该层旋钮上做调整,调完后重新测量,不能多层同时乱拧。每个旋钮都对应一个可观测量,本篇每一节都会把它挂出来:
1 | |
没有判据的调参就是碰运气:改完之后无法说明改善来自哪里,也无法判断该不该继续往同一个方向走。
JVM heap 调优
Logstash 运行在 JVM 上,heap 配置错误是生产环境最常见的性能隐患。
默认 heap 为 1 GB(-Xms1g -Xmx1g)。修改位于 config/jvm.options:
1 | |
-Xms 与 -Xmx 必须设为相同值,避免 JVM 在运行时动态扩展 heap——官方的说法是 resize 本身就是个代价很高的过程。
容器化部署里改 jvm.options 常常不方便,这时用 LS_JAVA_OPTS 环境变量:
1 | |
它的内容与 jvm.options 里的设置叠加,两处都出现同一个选项时 LS_JAVA_OPTS 覆盖文件里的值。
heap 上限的两条约束:
1 | |
至于 32 GB 那条 CompressedOops 阈值,在 Logstash 语境下用不上:推荐上限只有 8 GB,永远碰不到这个门槛。它是从 Elasticsearch 的调优经验里平移过来的约束,对 Logstash 不构成实际限制。
怀疑 heap 给得太小时,官方给了一个不需要任何工具的诊断:直接把 heap 翻倍,看性能是否改善。heap 太低会让 JVM 不停 GC,表现为 CPU 占用被无谓抬高。
off-heap:heap 与 PQ 之间的那笔账
-Xmx 只管一部分内存。操作系统、PQ 的 mmap page、direct memory、线程栈都在它之外,而本篇后半段要讲的 PQ 恰好就落在这一块里:
1 | |
三条需要记住的性质:
- direct memory 默认大小等于 heap。设了
-Xmx8g就意味着 JVM 默认还会再要 8 GB 的 direct memory 额度。官方建议考虑把-XX:MaxDirectMemorySize设成 heap 的一半,或者任何能容纳预期负载的值。 - 每条 PQ pipeline 至少要 head 和 tail 两个 page 常驻可访问,默认 page 是 64 MB,也就是每条 PQ pipeline 约 128 MB 起步。这是 per-pipeline 的固定成本,多管道场景下会线性累加。
- mmap 文件的大小无法设上界。这一块没有对应的 JVM 参数可以约束,只能在容量规划时预留。
官方给了整机估算公式:
1 | |
代入官方的例子:10 条 pipeline,每条 14 个线程(1 个 pipeline 线程 + 1 个 input 线程 + 12 个 worker),栈按 1 MB 算,heap 4 GB——原生内存 10 × (14 × 1MB + 128MB) = 1.4GB,direct memory 4 GB,heap 4 GB,合计约 9.4 GB。也就是说 -Xmx4g 这台机器实际要预留的是十来个 GB,不是 4 GB 加一点余量。
GC 算法:Logstash 近期版本默认使用 G1GC,是合理的默认选择。对于延迟敏感场景可以考虑 ZGC(-XX:+UseZGC),它的停顿时间更短但吞吐略低。GC 日志对于排查 heap 问题是必要的:
1 | |
heap 问题的判据是曲线形态,不是某个固定水位:已用堆与上限之间是否还留有余量、gc.collectors.old.collection_time_in_millis 的增长是平滑还是阶梯式上跳、hot threads 里有没有 GC 相关线程。已用堆长期贴着上限、老年代回收耗时快速累积,才是需要动手的信号。处置方式是先按上面那条"翻倍看是否改善"确认方向,再检查是否有 filter 在 event 里附加了大量不必要字段(每个 worker 持有一批 event 的引用,event 膨胀直接放大 heap 占用)。
pipeline.workers 与 pipeline.batch.size
这两个参数控制 filter 和 output 阶段的并发与批量。调它们之前先确认一件事:pipeline.workers 起的那组线程同时执行 filter 和 output,Logstash 里不存在独立的 output 线程池。一个 worker 在一次循环里读一个 batch、跑完全部 filter、再把结果交给 output,然后才回头取下一个 batch。所以下面所有"是 filter 慢还是 output 慢"的讨论,都发生在同一组线程上——output 卡住时 filter 不是"还在并行跑",而是跟着一起停。
默认值先摆出来,否则无从判断自己是在调大还是调小:
1 | |
pipeline.workers 的调优方向要按瓶颈是 CPU 还是 I/O 分开看,两种情况的上界完全不同:
1 | |
判据指标是 flow.worker_utilization:接近 100 说明 worker 已被占满,加 workers 有意义;明显低于 100 时瓶颈在别处,加 workers 只是多开几个闲着的线程。再往下要用插件级 worker_utilization 确认占满 worker 的是哪个插件,以及它是 CPU 型还是 I/O 型。
pipeline.batch.size 的调优方向:
1 | |
判据是批量的实际填充度。Logstash 在 node_stats 里暴露了 pipelines.<id>.batch.event_count,包含 average / p50 / p90 三组分时段统计(这项能力较新,9.4.0 起为技术预览;另有 pipeline.batch.metrics.sampling_mode 控制采样粒度)。读法:
1 | |
再叠加背压:低背压 + 满批 = 健康,可以考虑继续调大 batch 或 workers;高背压 + 满批 = 瓶颈在下游,该加 workers 或资源而不是加 batch;低背压 + 不满批 = 管道本来就没吃满,不需要调。
pipeline.batch.delay(默认 50ms)控制 worker 拿到至少一个 event 后等待凑满一批的超时。官方对它的判断是"rarely needs to be tuned",原因藏在一个乘积公式里:
1 | |
这不是一个可以忽略的量级。按下面那份配置代入(500 × 50ms),最坏情况下一条 event 要等 25 秒才进 filter。低速流量下这个最坏值是会被真正触及的——每次都差一点凑不满,每次都等满 50ms。
所以 batch.delay 基本不用动。真要动它,得先确认批量确实填不满(P50/P90 明显低于 batch.size),而且改完要盯住两个反向信号:PQ 场景看 queue.events 有没有涨,内存队列场景看 queue_push_duration_in_millis 有没有变长。这两个指标任一上升,说明管道开始向上游施加背压,整体性能是往下走的。
典型生产配置参考(仅作基准,需实测调整):
1 | |
持久队列磁盘调优
持久队列(PQ)把 event 落盘,提供 at-least-once 语义,但引入了磁盘 I/O 路径。磁盘成为瓶颈时的表现:events.queue_push_duration_in_millis 增长快,input 线程卡在往队列写这一步(hot threads 里能看到 [<id>]<inputname 的栈停在队列写入上)。
1 | |
写入侧和确认侧各有自己的 checkpoint 频率,对应两个不同的参数——这个区分在下面的取舍里是关键。
关键参数:
1 | |
queue.checkpoint.writes 是 PQ 调优里唯一一条真正的取舍,而它的代价方向很容易记反。官方在两处把话说得很直白。一处列在"PQ 解决不了的问题"里:“Data may be lost if an abnormal shutdown occurs before the checkpoint file has been committed”;另一处讲降低 checkpoint 频率时:“any data that has not been checkpointed, is lost”。
也就是说,已经写进 page 文件但还没进 checkpoint 的那批 event,在异常关停或硬件故障时是永久丢失,不是"还在 PQ 里等着重启后重处理"。调大这个值等于把"崩溃时最多永久丢多少条"这个数字一起调大。
两个极端官方都给了:
1 | |
SSD 与 HDD 的差异在 checkpoint 的随机写上最为明显:HDD 每次 fsync 耗时数毫秒到数十毫秒,直接拖慢 PQ 写入;SSD 的随机写延迟低一到两个数量级。生产环境 PQ 建议放 SSD。
只有 HDD 时,把 queue.checkpoint.writes 从 1024 调到 4096 确实能把 fsync 频率降到四分之一,但要把账算清楚:这个改动把崩溃时的永久丢失上限从 1024 条提到 4096 条。这笔交易值不值,取决于这条管道的数据能不能从源头重放——如果上游是 Kafka 或别的可重放的源,丢的这批还能再拉一次,代价可控;如果上游是 UDP syslog 这类无法重放的源,4096 条就是真的没了。
判据指标是 flow.queue_persisted_growth_events。官方的说法是它应该趋近于零;持续为正说明写入快于消费、队列在涨,负值代表队列正在收缩、积压在恢复。调完 checkpoint.writes 之后看这个值有没有从正数回落,比看磁盘 %util 更直接。
插件层调优
插件层的旋钮不在 logstash.yml 里,改动方式也不同:只改管道配置文件的话,可以靠 config.reload.automatic 自动重载,或者给进程发 SIGHUP 手动触发一次。前提是管道里的插件都支持 reload——input 和 output 插件常常持有 OS 资源,有些资源不重启进程就释放不掉,官方举的例子是 stdin input,它的存在会直接阻止整条管道重载。另有 pipeline.recoverable 影响重载失败后的可恢复行为。
filter 执行顺序对性能有直接影响:
1 | |
filter 顺序和插件选型的判据都是插件级 worker_millis_per_event:调整前后看同一个插件的单条 event 成本有没有下降,以及它在插件级 worker_utilization 里占的比例有没有缩小。
codec 选型:plain codec + dissect filter 的组合,在固定格式日志上比 json codec 或 grok 有更低的 CPU 开销。json codec 适合输入本身已经是结构化 JSON 的场景,无需 filter 再解析。避免在 json codec 解析后又用 grok 再次解析 message 字段——这是两次序列化/反序列化的叠加。
ES output 的 bulk 大小没有独立旋钮
这一节最需要说清的是一件"没有旋钮"的事:ES output 发出的 bulk 请求大小就是 pipeline.batch.size。官方对 pipeline.batch.size 的说明里写得很明确,ES output 对每个收到的 batch 尝试发一个 bulk 请求;调 pipeline.batch.size 就是在调 bulk 的大小。
配置里因此不需要、也不能出现 bulk 相关的参数:
1 | |
历史上的 flush_size 与 idle_flush_time 已经在 8.0.0 被标记 obsolete,后续版本的 changelog 明确写了 “Removed obsolete”,当前版本里已无踪迹。它们的后果比"参数不生效"重得多:Logstash 遇到不认识的插件参数会报 Unknown setting 'flush_size' for elasticsearch 并拒绝启动。抄一份旧配置进去,得到的不是一条警告,而是管道根本起不来。
output 侧真正可调的是连接池,与批量无关:
1 | |
这里还埋着一个比上面那两个更阴的反例。output plugin 上有一个 workers 参数,写进配置不报错、管道正常启动,但它在 logstash-core/lib/logstash/outputs/base.rb 里的声明是 config :workers, :deprecated => "This parameter will be ignored."——静默忽略,只在日志里留一条 deprecation。ES output 也没有 batch_size 这个配置项。
两类错误的危险程度正好相反于直觉:flush_size 让管道起不来,读者当场就知道要改;workers 让管道正常跑,读者以为 output 并发已经调过了,此后所有基于这个前提的调优推断全都建在空地上。静默无效比启动失败难查得多。
调优工作流
调优不应是随机调参,而是一个有方向的循环:
1 | |
常见陷阱:同时调多个参数,改善时不知道是哪个起效,恶化时不知道该回滚哪个。调优日志需要记录:参数名、旧值、新值、调整时间、测量结果。
模式提炼
1 | |
工程迁移表
| Logstash 调优概念 | Kafka 生态对应 | Flink 对应 | 通用 ETL 对应 |
|---|---|---|---|
-Xmx heap 上限 |
broker/consumer JVM heap | TaskManager heap | 任意 JVM 服务 heap |
| off-heap(PQ mmap + direct memory) | page cache + 堆外缓冲 | Managed / Network memory | 进程 RSS 中的非堆部分 |
pipeline.batch.size(同时决定 ES bulk 大小) |
consumer max.poll.records |
operator buffer 大小 | Extract/Load 批量 |
pipeline.workers |
consumer 线程数 / 分区数 | 算子并行度 | Transform 并发数 |
PQ checkpoint.writes / checkpoint.acks |
producer acks + fsync 频率 |
checkpoint 间隔 | Staging 提交批量 |
| filter 顺序(cheap first) | UDF 执行顺序优化 | 算子链顺序 | Transform 步骤排序 |
| dissect 替代 grok | 结构化日志直接反序列化 | 使用 RowData 替代 POJO | 固定格式直接解析 |
ES output pool_max / pool_max_per_route |
producer 连接数 / max.in.flight |
Sink 客户端连接池 | Load 阶段连接池 |
常见误解
误解一:“把 heap 调大一定能提速”。heap 过大会让 GC 的单次停顿时间变长(G1GC 需要扫描更大的 old gen),在延迟敏感场景下适得其反。调 heap 的目标是消除频繁 GC,而不是越大越好。先看 heap_used_percent 和 gc.collection_time_in_millis 增长速率,再决定是否需要扩 heap。
误解二:“不区分 CPU-bound 与 IO-bound,照 CPU 核数设 pipeline.workers”。核数只是 CPU 密集场景下的上界。官方明确写了 workers 可以超过核数——“may be set higher than the number of CPU cores since outputs often spend idle time in I/O wait conditions”,另一处更进一步:“Good results can even be found increasing this number past the number of available processors”。理由正是 I/O wait:这些线程大部分时间在等下游响应,超配 workers 才能把等待的时间利用起来。此时真正的上界是 heap 能承载的 workers × batch.size,而不是核数。反过来,如果瓶颈是 grok 这类 CPU 密集操作,加到核数附近就饱和了,再加只会让上下文切换吃掉收益。先用 worker_utilization 和热点线程的栈位置判断是哪一类,再定上界。
误解三:“PQ 里的数据崩溃后总能重来一遍”。要分写入侧和 ack 侧两头看,它们由两个不同的参数管:
1 | |
所以 PQ 的语义是 at-least-once 而不是 exactly-once(重复来自 ack 侧),同时它也不是"零丢失"(丢失来自写入侧)。要做到 exactly-once,还需要下游的幂等写配合(ES 用 document_id 去重,Kafka 用事务)。
误解四:“output plugin 上的 workers 参数能提高 output 并发”。这个参数确实存在、写进配置也不会报错,但它的声明带着 :deprecated => "This parameter will be ignored."——被完全忽略,只在日志里留一条 deprecation。output 侧的并发要靠连接池参数,ES output 是 pool_max 和 pool_max_per_route。这类静默无效的参数比拼错参数名更难发现:拼错会让管道起不来,静默无效会让人以为已经调过了。
误解五:“filter 顺序不影响性能,只影响正确性”。顺序直接影响性能。如果把高 CPU 的 grok 放在最前面,每个 event 都要跑一遍 grok,哪怕后续 filter 会丢弃大部分 event。把廉价的条件判断(if [type] == "app")放在 grok 之前,能让大量不需要 grok 的 event 直接跳过,显著降低 CPU 负载。
练习
-
构造一条带 PQ 的 Logstash 管道,向它注入高速事件流,用
iostat -x 1观察磁盘%util和await,对比 HDD 和 SSD(或用dm-delay模拟 HDD 延迟)时events.queue_push_duration_in_millis的差异。再把queue.checkpoint.writes从 1024 调到 4096,重复测量,同时记录flow.queue_persisted_growth_events的变化。最后把这次调优的代价写成一句话:这条管道在崩溃时的永久丢失上限从多少条变成了多少条,以及上游能不能重放。 -
设计两条管道跑同一份 workers 扫描实验:一条是 CPU 密集的(重 grok,output 用
stdout { codec => dots }),一条是 I/O 密集的(filter 极简,output 指向一个人为加了 200ms 延迟的http端点)。两条都把pipeline.workers从 1 逐步加到 2×核数,各值采集 30 秒的吞吐与flow.worker_utilization,画出两条曲线。确认 CPU 密集那条在核数附近就转平或下滑,而 I/O 密集那条可以超过核数继续涨,并用热点线程的栈位置解释差异。 -
把
pipeline.batch.size设为 500、pipeline.batch.delay保持 50ms,然后用极低速率(比如每 2 秒一条)注入事件,测量单条 event 从注入到出现在 output 的实际延迟。对照batch.delay × batch.size算出的最坏值 25 秒,看实测落在什么位置;再把batch.size降到 125 重测一次,量化这个乘积对低速流量的影响。
系列导航
参考资料
- Logstash 调优指南:https://www.elastic.co/guide/en/logstash/current/tuning-logstash.html(workers 可超核数的原文表述、
batch.delay × batch.size的 worst-case 公式、“Optimizing batch sizes” 一节的 P50/P90 判读表、ES output 每个 batch 一个 bulk 请求) - Logstash 性能问题排查:https://www.elastic.co/guide/en/logstash/current/performance-troubleshooting.html(瓶颈分层排查流程)
- Logstash 持久队列配置:https://www.elastic.co/guide/en/logstash/current/persistent-queues.html(
page_capacity、checkpoint.writes/checkpoint.acks与两个极端值、未 checkpoint 数据"is lost"的原文、queue.drain) - Logstash JVM 配置:https://www.elastic.co/guide/en/logstash/current/jvm-settings.html(4–8GB 建议区间、物理内存 50–75%、off-heap 三类内存的构成表、整机内存估算公式、
LS_JAVA_OPTS与jvm.options的叠加与覆盖规则) - Logstash 配置重载:https://www.elastic.co/guide/en/logstash/current/reloading-config.html(
config.reload.automatic、SIGHUP,以及哪些插件会阻止 reload) - G1GC 调优指南(Oracle):https://docs.oracle.com/en/java/javase/17/gctuning/garbage-first-g1-garbage-collector1.html
- Logstash dissect filter 文档:https://www.elastic.co/guide/en/logstash/current/plugins-filters-dissect.html(与 grok 的性能对比说明)
- Logstash elasticsearch output 文档:https://www.elastic.co/guide/en/logstash/current/plugins-outputs-elasticsearch.html(
pool_max/pool_max_per_route连接池参数) - 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.")
