上一篇(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
14
Memory Queue(queue.type: memory)

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

RAM 里的有界环形缓冲区
容量 = batch.size × workers(不可独立配置)
进程被强杀 → 全部丢失

Persistent Queue(queue.type: persisted)

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

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

内存队列是一块有界的 RAM 缓冲区,速度最快,但进程被强杀或崩溃时,队列里尚未被 output 确认的 event 全部消失。

它的"有界"到底有多大,比多数人预期的小得多。内存队列的容量不是一个独立参数,而是由 pipeline.batch.size × pipeline.workers 算出来的,默认 125 × CPU 核数,8 核机器上约 1000 条。底层实现是 ArrayBlockingQueue,内部一个循环数组,所以"有界环形缓冲区"这个说法准确,只是缓冲深度只有一千条量级。PQ 那一侧的 queue.max_bytes 默认是 1024mb,两者不在同一个量级上,后面性能权衡表里那些取舍的参照物就是这个差距。

持久队列把 event 写进磁盘上的 page 文件,并维护一个 checkpoint 文件记录"哪些 event 已经被 output 成功处理"。崩溃重启后,Logstash 从 checkpoint 读取进度,重放尚未被确认的 event。

版本说明

queue.type 的默认值至今仍是 memory,6.x 到 9.x 一路没变过,开 PQ 永远是一个显式动作。真正需要对着当前版本文档核一遍的是队列参数清单本身:queue.compression 是较新加入的(取值 none / speed / balanced / size / disabled,PQ 文档上标着 stack: ga 9.2),queue.checkpoint.retry 控制 checkpoint 写失败后是否重试(默认 true)。老版本上写这些参数会被当成 unknown setting 直接拒绝启动。

持久队列的内部结构

PQ 在磁盘上的布局如下:

1
2
3
4
5
6
7
8
path.queue/                 # 默认 path.data/queue
├── .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,默认 1024。这个值越小,checkpoint 越频繁,崩溃时需要重放的 event 越少,但磁盘 I/O 开销越高。写入侧之外还有对称的确认侧:queue.checkpoint.acks 控制每 ACK 多少条 event 做一次 checkpoint,默认同样是 1024。两个参数各管一条路径上的 fsync 节奏,调优时要一起看。

PQ 这一组参数的默认值列在一起更容易建立量感:

参数 默认值 含义
queue.type memory 队列类型,开 PQ 要显式改成 persisted
path.queue path.data/queue PQ 文件目录
queue.max_bytes 1024mb PQ 占用磁盘的总上限
queue.max_events 0 队列内 event 数上限,0 表示不限
queue.page_capacity 64mb 单个 page 文件大小
queue.checkpoint.writes 1024 每写入多少条做一次 checkpoint
queue.checkpoint.acks 1024 每 ACK 多少条做一次 checkpoint
queue.checkpoint.retry true checkpoint 写失败时是否重试
1
2
3
4
5
6
7
# logstash.yml PQ 相关配置示例
queue.type: persisted
path.queue: /var/lib/logstash/queue # 不写则用 path.data/queue
queue.max_bytes: 4gb # 示例值,默认是 1024mb
queue.page_capacity: 64mb # 单个 page 文件大小
queue.checkpoint.writes: 1024 # 每写 1024 条做一次 checkpoint
queue.checkpoint.acks: 1024 # 每 ACK 1024 条做一次 checkpoint

退出时保住什么:queue.drain

崩溃是一条边界,正常退出是另一条,两条走的代码路径完全不同。

收到 SIGTERM(Ctrl-C、systemctl stop、容器停止)时,Logstash 先停掉 input,然后让 worker 继续跑。跑到什么时候停下,由 queue.drain 决定,默认 false

  • false:worker 把手上那批 event 处理完就退出,队列里剩下的不再消费。PQ 下这些 event 留在磁盘上,下次启动接着处理。
  • true:退出前把队列排空,所有已入队的 event 都处理完才结束进程。代价是关闭时长变成不确定的,队列有多深就得等多久。

这里有个容易被这个参数名带偏的实现细节:内存队列不受它影响,优雅关闭时总是排空的。Logstash 内部计算 drain 标记的条件是"queue.drain 为真,或者队列类型是 memory",两者取或。理由也直白:内存队列没有下次启动可以接着处理的落盘状态,不排空就等于丢弃。所以 queue.drain 实际是一个只对 PQ 生效的开关,它决定的是"PQ 里剩下的 event 现在处理完,还是留给下次启动"。

kill -9 和断电是另一条路径:进程被直接终止,没有任何清理机会,queue.drain 无从介入。内存队列里的 event 全部消失,PQ 里的靠 checkpoint 恢复。这也是前面"内存队列进程退出全部丢失"那句话需要限定的地方,它成立的前提是强杀或崩溃,不是优雅关闭。

优雅关闭本身还有卡住的可能:某个 output 一直阻塞,worker 收不了尾。默认的 pipeline.unsafe_shutdown: false 下 Logstash 会一直等,只在日志里反复报 The shutdown process appears to be stalled due to busy or blocked plugins。把它设成 true,Logstash 在连续几轮判定停滞之后强制退出,代价是在途 event 按 kill -9 的规则处理。

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
path.queue: /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 文件在磁盘上 写入路径是 org.logstash.ackedqueue.QueuePageMmapPageIOV2,page 文件用 mmap 写,序列化格式为 Logstash 自定义的 binary
checkpoint.head 文件 Checkpoint 对象记录 head page 序号和写偏移量
kill 后重启,event 重放 启动时读取 checkpoint,重放起点是 firstUnackedPageNum + firstUnackedSeqNum
已输出的 event 不重复 org.logstash.ackedqueue.Queue 的 ACK 机制:output 成功后回调 batch.close()
page 文件最终消失 page 内所有 event 全部 ACK → PageIO.purge() 删除文件

第三行那个字段名值得多看一眼,因为它的方向和直觉相反。Checkpoint 里没有"已经 ACK 到哪"这种字段,它记的是第一个还没被 ACK 的位置,全部字段就是 pageNum / firstUnackedPageNum / firstUnackedSeqNum / minSeqNum / elementCount 五个。重放从"第一个未确认"往后走,而不是从"最后一个已确认"往后走。两种记法在正常情况下等价,但前者在 checkpoint 本身落后于实际 ACK 进度时是安全的(重放多几条),后者会漏。这就是 at-least-once 而不是 exactly-once 在数据结构层面的体现。

内存队列 vs 持久队列:性能权衡

PQ 带来可靠性的代价是 I/O 延迟。每条 event 写入 PQ 时需要序列化并写磁盘。

1
2
3
4
5
6
7
8
9
10
性能权衡对比:

维度 内存队列 持久队列
──────────────────────────────────────────────────
写入延迟 内存操作 序列化 + 磁盘写入,随 fsync 频率放大
崩溃恢复 无法恢复 从 checkpoint 重放
容量上限 batch.size × workers queue.max_bytes(默认 1024mb)
吞吐峰值 更高 取决于 fsync 频率与磁盘介质
JVM 堆压力 高(event 留在堆) 低(序列化后落磁盘)
适用场景 允许丢数据的场景 生产环境、数据不可丢失

“随 fsync 频率放大"和"取决于介质"这两格不能压缩成"中等"或"略低”,因为 PQ 的开销主体不是"写磁盘"这个笼统说法,而是 fsync 的次数,而次数由两个参数直接决定:queue.checkpoint.writesqueue.checkpoint.acks。默认每 1024 条各做一次,摊到单条 event 上几乎看不见;把它们设成 1(每条都 fsync),fsync 次数直接乘 1024:这是量级层面的差别,"中等"和"略高"这种刻度描述不了。

第二个变量是介质。同样一份配置跑在本地 NVMe、机械盘和网络块存储上,fsync 延迟差一到两个数量级,PQ 的吞吐损失可以从个位数百分比一路到数倍。第三个变量是 page 的序列化开销,它随单条 event 的大小走。

三个变量叠起来的结果是:PQ 相对内存队列的性能代价没有一个通用比例可引用,要给自己的环境一个数字,只能在目标介质上按目标 event 形态跑一次对照实验。

可靠性放在 Logstash 内还是外:PQ vs 前置 Kafka

PQ 的边界前面已经划清了:它保住的是"进入队列之后",进入之前那段以及源本身的重放能力它管不了。而让源具备重放能力这件事,在工程上通常只有一个答案——在 Logstash 前面放一个 Kafka(或同类的持久化消息系统)。这是可靠性设计的第一个岔路口:把缓冲和重放放在 Logstash 进程内,还是放在它外面。

1
2
3
4
5
6
7
8
9
方案 A:source → Logstash(PQ) → ES
缓冲在 Logstash 本地磁盘,一份组件、一套运维
重放边界 = PQ 里还没 ACK 的部分
Logstash 实例本身是单点:机器挂了,PQ 文件跟着不可用

方案 B:source → Kafka → Logstash(memory) → ES
缓冲在 Kafka,按 topic retention 计算,天级容量是常态
重放边界 = Kafka 的 offset,可以任意回退重跑
Logstash 变成无状态消费者,可以横向扩、可以整机替换

两者的取舍不在性能而在边界位置:

维度 PQ 前置 Kafka
缓冲深度 本机磁盘,queue.max_bytes 量级通常在 GB topic retention,量级通常在 TB / 天
重放能力 只能重放未 ACK 的部分,已 ACK 的无法回退 可以把 offset 回退到 retention 内任意位置重跑
Logstash 有状态性 有状态:PQ 目录与实例绑定,.lock 排他 无状态:实例可随意增减、替换
多消费方 无:event 出队即消失 有:多个 consumer group 各自消费同一份数据
运维成本 一个参数 一套 Kafka 集群

判断顺序可以简化成两问:需要的缓冲深度是否超过单机磁盘能给的量级,以及是否需要"把历史数据重新跑一遍"的能力。两个都是"否",PQ 足够,多一套 Kafka 是净成本;任何一个是"是",PQ 都补不上,该上 Kafka。

两者不互斥。Kafka → Logstash(PQ) → ES 是常见组合,此时 PQ 的作用退化成一层薄的本地兜底,让 Logstash 在 ES 短暂不可用时不必立刻把背压顶回 Kafka。

模式提炼

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 一定比内存队列慢很多”。默认配置下 checkpoint 每 1024 条才触发一次 fsync,其余写入落在 OS page cache 上,正常路径的 I/O 延迟增加有限。PQ 把 event 从 JVM 堆序列化到磁盘,还顺带把这部分对象移出了堆,在堆吃紧的高吞吐场景下有时比内存队列更稳定。真正会让 PQ 掉一个量级的是把 checkpoint 参数调到极小,或者把 PQ 目录放在网络盘上。

误解三:“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 的磁盘上限(默认 1024mb),它能做到的是在下游持续慢时给 Logstash 更多缓冲空间,避免过早触发背压。但它不能替代源的重放能力:一旦 PQ 写满,背压还是会传回 input,新 event 仍然面临被丢弃(若 input 是 UDP)或被阻塞(若 input 是 Beats/Kafka)的问题。要的是"能把历史数据重跑一遍"这种能力,queue.max_bytes 调到多大都给不了,那是前置 Kafka 的职责。

误解五:“优雅关闭不会丢数据,所以内存队列在受控重启时是安全的”。前半句成立,后半句的推理断了。优雅关闭时内存队列确实会排空,但"优雅"的前提是进程收到 SIGTERM 并且能顺利收尾。output 阻塞时关闭会停滞,容器编排器等不到超时就补一个 SIGKILL,pipeline.unsafe_shutdown: true 也会在判定停滞后主动强制退出,这几种情况下内存队列同样全丢。受控重启的安全性取决于关闭流程能不能走完,不取决于它名义上是不是"正常退出"。

练习

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

  2. 用内存队列做同样的实验(queue.type: memory,同样 kill -9),确认崩溃时的 event 丢失行为。然后改成用 Ctrl-C(SIGTERM)退出再跑一次,对比两种退出方式下内存队列的差别,并说明为什么 queue.drain 在这两次实验里都没有参与。

  3. 思考题:一个"Beats → Logstash(PQ)→ Elasticsearch"的管道,Logstash 崩溃后重启。分析以下三种 event 各自的命运:(a) 已经被 ES output 成功写入并 ACK 的;(b) 在 PQ 里等待被 worker 消费的;© Beats 已发送但 Logstash 崩溃前还没写入 PQ 的。Beats 的 ack 机制在第 3 种情况下如何介入?

系列导航

序号 主题
00 导读:核心对象是 event,骨架是三段管道
01 架构:JRuby、JVM 与 pipeline 的运行形态
02 event 模型:@timestamp、@metadata 与字段引用
03 codec:字节流与 event 的边界转换
04 input 插件:拉取、监听与 Beats 接入
05 Grok 的本质:命名正则加预定义 pattern
06 dissect 与结构化 filter:放弃回溯换吞吐
07 常用 filter 组合:mutate、date、geoip 与条件
08 output 插件:Elasticsearch output 与批量写入
09 pipeline 执行模型:worker、batch 与背压
10 内存队列 vs 持久队列:可靠性的分界线(本篇)
11 死信队列(DLQ):无法处理的 event 去哪
12 Multiple Pipelines 与 pipeline-to-pipeline
13 监控:Node Stats API、hot threads 与瓶颈定位
14 性能调优:JVM heap、批处理与持久队列磁盘
15 Logstash vs Beats vs Ingest Pipeline:该用谁
16 Logstash vs Fluentd vs Vector:日志管道的三种取舍
17 Logstash 的演进与 Elastic Agent 的冲击

参考资料