深入 Logstash 10 - 内存队列 vs 持久队列:可靠性的分界线
上一篇(09)拆解了 pipeline 执行模型:worker 并发、batch 凑批、背压自然传导。这一篇进入队列本身:内存队列和持久队列的结构差异,崩溃时各自保住了什么、保不住什么,以及 at-least-once 语义的真实边界。
核心问题:开启持久队列(PQ)之后,Logstash 崩溃重启能恢复哪些 event,又有哪些情况它依然无能为力?
两种队列的基本结构
Logstash pipeline 内部在 input 和 filter/output 之间有一个队列,起到解耦和缓冲的作用。队列类型由 logstash.yml(或 pipelines.yml)里的 queue.type 参数控制。
1 | |
内存队列是一块有界的 RAM 缓冲区,速度最快,但进程一旦退出(正常退出或崩溃),队列里尚未被 output 确认的 event 全部消失。
持久队列把 event 写进磁盘上的 page 文件,并维护一个 checkpoint 文件记录"哪些 event 已经被 output 成功处理"。崩溃重启后,Logstash 从 checkpoint 读取进度,重放尚未被确认的 event。
版本说明
queue.type 的默认值在不同版本之间有过变化,使用前应以当前版本的官方文档为准:
- Logstash 6.x 及更早:默认
memory - Logstash 7.x:默认
memory,官方文档推荐在生产环境评估 PQ - Logstash 8.x:默认仍为
memory(截至 8.x 主线,官方未将 PQ 设为默认)
不同发行版(如 Elastic Cloud 托管的 Logstash)可能有不同的默认配置,应以具体版本的 Release Notes 为准,不应假设"Logstash 一直默认内存队列"或"8.x 已经默认 PQ"。
持久队列的内部结构
PQ 在磁盘上的布局如下:
1 | |
写入流程:input 把 event 序列化后追加到 head page。当 head page 达到 queue.page_capacity(默认 64MB)时,它变成 tail page,并打开新的 head page。
读取流程:worker 从最旧的 tail page(或 head page,如果队列只有一个 page)按顺序读取 event,组成一批返回给 pipeline。
ACK 流程:worker 完成一批 event 的 output 之后,向队列发送 ACK,告知"这些序号的 event 已经成功处理"。PQ 更新对应 page 的 checkpoint。当一个 page 上所有的 event 都被 ACK,该 page 文件被删除,磁盘空间释放。
1 | |
checkpoint 与崩溃恢复
checkpoint 文件是 PQ 可靠性的核心。它记录了"截至某个时刻,哪些 event 序号已经被 output 确认"。
崩溃重启时,Logstash 读取所有 checkpoint 文件,找到最后一次成功 ACK 的位置,从该位置之后开始重放。这保证了"进入 PQ 之后"的 event 在崩溃后不会永久丢失。
queue.checkpoint.writes 参数控制每写入多少个 event 做一次 checkpoint(默认值在不同版本间有差异,需查阅当前版本文档)。这个值越小,checkpoint 越频繁,崩溃时需要重放的 event 越少,但磁盘 I/O 开销越高。
1 | |
at-least-once 的真实边界
PQ 提供的语义是 at-least-once,即"同一个 event 可能被处理多次,但不会永久丢失"。"多次"来自崩溃时最后一个未完整 ACK 的 batch:重启后这批 event 会被重放,而它们可能已经被 output 部分写入了下游。
更重要的是 at-least-once 的覆盖范围边界:
1 | |
理解第 1 条的关键:input 把 event 写入 PQ 不是原子操作,写入过程本身如果被中断,该 event 可能还没落盘。这段"input 刚读到但还没落入 PQ"的窗口期,PQ 无法覆盖。
第 2 条说明了 PQ 的先决条件:源可重放。如果数据源能重放(Kafka 保留 offset、Beats 支持 ACK 重发、文件 input 记录读取位点),崩溃后重放 PQ 里的 event 并不会导致从源头丢失数据。如果数据源不可重放(UDP、无状态的 TCP 一次性发送),进入 PQ 之前这段已经丢了的数据,PQ 不能补回。
1 | |
实验:观察 PQ 落盘行为
用如下最小配置开启 PQ:
1 | |
1 | |
启动后输入 10 行文本,然后在 filter sleep 期间强制 kill -9 Logstash 进程:
1 | |
查看 /tmp/logstash-pq-test/ 目录,可以看到 page.0(或 page.N)和 checkpoint.* 文件。重启 Logstash,观察已经进入 PQ 的 event 是否被重放输出——已经输出过的(被 ACK 的)不会重复,未被 ACK 的会重新出现。
对应内部对象
| 实验现象 | 对应内部结构 |
|---|---|
| page.0 文件在磁盘上 | PQWriter 写入的 page 文件,序列化格式为 Logstash 自定义的 binary |
| checkpoint.head 文件 | Checkpoint 对象记录 head page 序号和写偏移量 |
| kill 后重启,event 重放 | 启动时读取 checkpoint,从 min_acked_count 之后重放 |
| 已输出的 event 不重复 | AckedQueue 的 ACK 机制:output 成功后回调 batch.close() |
| page 文件最终消失 | page 内所有 event 全部 ACK → PageIO.purge() 删除文件 |
内存队列 vs 持久队列:性能权衡
PQ 带来可靠性的代价是 I/O 延迟。每条 event 写入 PQ 时需要序列化并写磁盘(取决于 checkpoint 频率和 OS page cache 刷盘策略)。
1 | |
一个反直觉的好处:PQ 把 event 序列化后存磁盘,event 对象从 JVM 堆上释放,对需要高吞吐且 JVM 堆有限的场景,PQ 有时反而能降低 GC 压力,间接提升稳定性。
模式提炼
1 | |
这套模式和 Kafka Consumer 的 offset 提交是同构的:Kafka consumer 读取消息但不提交 offset,崩溃后从上次 offset 重放——这就是 Kafka 侧的"进入 broker 之后"的 at-least-once。PQ 的 checkpoint 是同一个思路在 Logstash 内部的实现。
工程迁移表
| Logstash 概念 | Kafka 对应 | Flink 对应 | 通用 ETL 对应 |
|---|---|---|---|
| queue.type: memory | 无持久化的本地缓冲 | 无 checkpoint 的流处理 | 纯内存 staging |
| queue.type: persisted | Kafka broker(本身就是持久日志) | RocksDB state backend | 落盘 staging 区 |
| checkpoint 文件 | consumer offset(提交到 __consumer_offsets) |
Flink checkpoint(存 HDFS/S3) | 水位线 / 进度表 |
| ACK after output success | commitSync() / commitAsync() | checkpoint barrier 对齐 | 成功写目标后更新水位 |
| at-least-once 边界 | 消息进 broker 之后 | 事件进入 Flink 处理之后 | 数据进入 staging 之后 |
| queue.max_bytes | retention.bytes / log.retention.bytes | state.backend.fs.memory-threshold | staging 磁盘配额 |
| page 文件被 purge | log segment 被 cleanup 删除 | completed checkpoint 被 discard | staging 区清理 |
常见误解
误解一:“开了 PQ 就不会丢数据”。PQ 只保证"进入队列之后"的 at-least-once。如果 input 从不可重放的源(UDP、某些脚本化的 stdin 管道)读取数据,数据在进入 PQ 前已丢失,PQ 无法补救。at-least-once 是有前提条件的,不能省略"源可重放"这个条件。
误解二:“PQ 一定比内存队列慢很多”。PQ 的写入走磁盘,OS 的 page cache 会缓冲大部分写操作,实际上在正常运行路径下,I/O 延迟增加有限。更重要的是,PQ 把 event 从 JVM 堆序列化到磁盘,可以降低 GC 压力,在高吞吐场景下有时反而比内存队列更稳定。
误解三:“at-least-once 意味着下游会收到重复数据,这是 Logstash 的问题”。重复是 at-least-once 语义的一部分,不是 Logstash 特有的缺陷。Kafka 的 at-least-once consumer 同样会在崩溃后重放未提交 offset 的消息。处理重复需要在目标端做幂等写(如 ES 的 document_id 去重)或者在整条链路上实现 exactly-once(代价更高),这是架构设计选择,不是 Logstash 的 bug。
误解四:“queue.max_bytes 越大越安全”。queue.max_bytes 是 PQ 的磁盘上限,它能做到的是在下游持续慢时给 Logstash 更多缓冲空间,避免过早触发背压。但它不能替代源的重放能力:一旦 PQ 写满,背压还是会传回 input,新 event 仍然面临被丢弃(若 input 是 UDP)或被阻塞(若 input 是 Beats/Kafka)的问题。
练习
-
按照本文实验步骤,在本地用
kill -9模拟崩溃,验证 PQ 的重放行为。记录哪些 event 被重放(未 ACK 的),哪些没有重放(已 ACK 的)。把观察结果和 checkpoint 文件的修改时间对照。 -
用内存队列做同样的实验(
queue.type: memory,同样kill -9),对比重启后的输出,验证内存队列在崩溃时的 event 丢失行为。 -
思考题:一个"Beats → Logstash(PQ)→ Elasticsearch"的管道,Logstash 崩溃后重启。分析以下三种 event 各自的命运:(a) 已经被 ES output 成功写入并 ACK 的;(b) 在 PQ 里等待被 worker 消费的;© Beats 已发送但 Logstash 崩溃前还没写入 PQ 的。Beats 的 ack 机制在第 3 种情况下如何介入?
系列导航
| 序号 | 主题 | 状态 |
|---|---|---|
| 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 | 演进、生态与对比 |
参考资料
- Logstash 持久队列官方文档:https://www.elastic.co/guide/en/logstash/current/persistent-queues.html(PQ 结构、配置参数、at-least-once 语义说明)
- Logstash queue 配置参数:https://www.elastic.co/guide/en/logstash/current/logstash-settings-file.html(queue.type、queue.max_bytes、queue.checkpoint.writes 等)
- Logstash 源码仓库:https://github.com/elastic/logstash(
logstash-core/src/main/java/org/logstash/ackedqueue/— AckedQueue、Checkpoint、Page 实现) - Beats ACK 机制文档:https://www.elastic.co/guide/en/beats/filebeat/current/how-filebeat-works.html(Beats 侧的 at-least-once 与 Logstash PQ 的配合)
- Kafka Consumer offset 文档:https://kafka.apache.org/documentation/#consumerconfigs(enable.auto.commit、auto.commit.interval.ms,与 PQ checkpoint 的同构对比)
- Elasticsearch 幂等写文档:https://www.elastic.co/guide/en/elasticsearch/reference/current/docs-index_.html(document_id 去重,处理 at-least-once 重放产生的重复)
