上一篇(09)拆解了 pipeline 执行模型:worker 并发、batch 凑批、背压自然传导。这一篇进入队列本身:内存队列和持久队列的结构差异,崩溃时各自保住了什么、保不住什么,以及 at-least-once 语义的真实边界。

核心问题:开启持久队列(PQ)之后,Logstash 崩溃重启能恢复哪些 event,又有哪些情况它依然无能为力?

两种队列的基本结构

Logstash pipeline 内部在 input 和 filter/output 之间有一个队列,起到解耦和缓冲的作用。队列类型由 logstash.yml(或 pipelines.yml)里的 queue.type 参数控制。

1
2
3
4
5
6
7
8
9
10
11
12
13
Memory Queue(queue.type: memory)

input thread ──▶ [ event | event | event | event | event ] ──▶ worker

RAM 里的有界环形缓冲区
进程退出 → 全部丢失

Persistent Queue(queue.type: persisted)

input thread ──▶ [ head page (写) | tail page | tail page ] ──▶ worker

磁盘 page 文件 + checkpoint 文件
进程崩溃重启 → 从 checkpoint 恢复

内存队列是一块有界的 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
2
3
4
5
6
7
8
queue.path/
├── .lock # 进程锁,防止多实例同时写
├── checkpoint.head # 当前写入位置(head page 序号 + 写偏移)
├── checkpoint.X # 每个 tail page 对应一个 checkpoint
├── page.0 # tail page(已写入,等待读取和 ACK)
├── page.1
├── ...
└── page.N # head page(当前正在写入的 page)

写入流程: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
2
3
4
5
6
写入顺序:
input → serialize → append to head page → (page full) → head → tail, new head

读取 + ACK 顺序:
worker ← batch ← tail page[0]
output success → ACK(batch) → update checkpoint → (page fully ACKed) → delete page file

checkpoint 与崩溃恢复

checkpoint 文件是 PQ 可靠性的核心。它记录了"截至某个时刻,哪些 event 序号已经被 output 确认"。

崩溃重启时,Logstash 读取所有 checkpoint 文件,找到最后一次成功 ACK 的位置,从该位置之后开始重放。这保证了"进入 PQ 之后"的 event 在崩溃后不会永久丢失。

queue.checkpoint.writes 参数控制每写入多少个 event 做一次 checkpoint(默认值在不同版本间有差异,需查阅当前版本文档)。这个值越小,checkpoint 越频繁,崩溃时需要重放的 event 越少,但磁盘 I/O 开销越高。

1
2
3
4
5
6
# logstash.yml PQ 相关配置示例
queue.type: persisted
queue.path: /var/lib/logstash/queue
queue.max_bytes: 4gb # PQ 总磁盘上限
queue.page_capacity: 64mb # 单个 page 文件大小
queue.checkpoint.writes: 1024 # 每写 1024 条做一次 checkpoint

at-least-once 的真实边界

PQ 提供的语义是 at-least-once,即"同一个 event 可能被处理多次,但不会永久丢失"。"多次"来自崩溃时最后一个未完整 ACK 的 batch:重启后这批 event 会被重放,而它们可能已经被 output 部分写入了下游。

更重要的是 at-least-once 的覆盖范围边界:

1
2
3
4
5
6
7
8
9
覆盖范围(PQ 能保障的):
event 进入 PQ 之后,在 output 成功 ACK 之前发生崩溃
→ 重启后重放,不永久丢失

不覆盖的范围(PQ 无能为力的):
1. event 还没进入 PQ 时崩溃
(input 从数据源读取了,但还没写入 queue.write())
2. 数据源本身不可重放(UDP、某些 stdin 场景)
3. 下游 output 不幂等(重放导致重复写入,产生副作用)

理解第 1 条的关键:input 把 event 写入 PQ 不是原子操作,写入过程本身如果被中断,该 event 可能还没落盘。这段"input 刚读到但还没落入 PQ"的窗口期,PQ 无法覆盖。

第 2 条说明了 PQ 的先决条件:源可重放。如果数据源能重放(Kafka 保留 offset、Beats 支持 ACK 重发、文件 input 记录读取位点),崩溃后重放 PQ 里的 event 并不会导致从源头丢失数据。如果数据源不可重放(UDP、无状态的 TCP 一次性发送),进入 PQ 之前这段已经丢了的数据,PQ 不能补回。

1
2
3
4
5
6
7
8
9
10
可靠性等级对比:

数据源 + 队列类型 → 故障时的表现
──────────────────────────────────────────────────────────
Kafka offset + memory → Logstash 崩溃后 Kafka offset 未提交,重启重放(Kafka 兜底,Logstash 队列不保证)
Kafka offset + PQ → 两层保障,at-least-once 端到端更稳健
Beats ack + memory → Logstash 崩溃后 Beats 重发,内存队列里的 event 丢失但 Beats 会补
Beats ack + PQ → at-least-once 端到端,PQ 和 Beats ack 双重兜底
UDP + memory → 崩溃丢失,无法恢复
UDP + PQ → 进入 PQ 前的 event 丢失,进入 PQ 后的可恢复,但 UDP 本身没有重传

实验:观察 PQ 落盘行为

用如下最小配置开启 PQ:

1
2
3
4
# logstash.yml
queue.type: persisted
queue.path: /tmp/logstash-pq-test
queue.max_bytes: 256mb
1
2
3
4
5
6
7
8
9
10
# pipeline10.conf
input {
stdin { codec => line }
}
filter {
sleep { time => "2" every => 1 } # 故意让 worker 慢下来,让 event 在队列里停留
}
output {
stdout { codec => rubydebug }
}

启动后输入 10 行文本,然后在 filter sleep 期间强制 kill -9 Logstash 进程:

1
2
3
bin/logstash -f pipeline10.conf
# 输入 10 行后,在另一个终端:
kill -9 $(pgrep -f logstash)

查看 /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
2
3
4
5
6
7
8
9
10
性能权衡对比:

维度 内存队列 持久队列
──────────────────────────────────────────
写入延迟 极低(内存操作) 中等(磁盘序列化)
崩溃恢复 无法恢复 从 checkpoint 重放
磁盘占用 无 queue.max_bytes 上限
吞吐峰值 更高 略低(受磁盘 I/O 限制)
JVM 堆压力 高(event 留在堆) 低(序列化后落磁盘)
适用场景 允许丢数据的场景 生产环境、数据不可丢失

一个反直觉的好处:PQ 把 event 序列化后存磁盘,event 对象从 JVM 堆上释放,对需要高吞吐且 JVM 堆有限的场景,PQ 有时反而能降低 GC 压力,间接提升稳定性。

模式提炼

1
2
3
4
5
6
7
模式:落盘检查点 + ACK 确认 = at-least-once 队列

- 把队列状态写入持久存储(磁盘),不只留在内存
- 消费者明确 ACK"已处理完",而不是"已读取"
- checkpoint 记录最后一次 ACK 的位置,崩溃后从这里重放
- at-least-once 的边界:进入持久存储之后,ACK 之前的窗口
- 端到端可靠性需要:源可重放 + 队列持久化 + 目标幂等写

这套模式和 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)的问题。

练习

  1. 按照本文实验步骤,在本地用 kill -9 模拟崩溃,验证 PQ 的重放行为。记录哪些 event 被重放(未 ACK 的),哪些没有重放(已 ACK 的)。把观察结果和 checkpoint 文件的修改时间对照。

  2. 用内存队列做同样的实验(queue.type: memory,同样 kill -9),对比重启后的输出,验证内存队列在崩溃时的 event 丢失行为。

  3. 思考题:一个"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 演进、生态与对比

参考资料