深入 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 全部消失。
它的"有界"到底有多大,比多数人预期的小得多。内存队列的容量不是一个独立参数,而是由 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 | |
写入流程: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,默认 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 | |
退出时保住什么: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 | |
第 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 文件在磁盘上 | 写入路径是 org.logstash.ackedqueue.Queue → Page → MmapPageIOV2,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 | |
“随 fsync 频率放大"和"取决于介质"这两格不能压缩成"中等"或"略低”,因为 PQ 的开销主体不是"写磁盘"这个笼统说法,而是 fsync 的次数,而次数由两个参数直接决定:queue.checkpoint.writes 和 queue.checkpoint.acks。默认每 1024 条各做一次,摊到单条 event 上几乎看不见;把它们设成 1(每条都 fsync),fsync 次数直接乘 1024:这是量级层面的差别,"中等"和"略高"这种刻度描述不了。
第二个变量是介质。同样一份配置跑在本地 NVMe、机械盘和网络块存储上,fsync 延迟差一到两个数量级,PQ 的吞吐损失可以从个位数百分比一路到数倍。第三个变量是 page 的序列化开销,它随单条 event 的大小走。
三个变量叠起来的结果是:PQ 相对内存队列的性能代价没有一个通用比例可引用,要给自己的环境一个数字,只能在目标介质上按目标 event 形态跑一次对照实验。
可靠性放在 Logstash 内还是外:PQ vs 前置 Kafka
PQ 的边界前面已经划清了:它保住的是"进入队列之后",进入之前那段以及源本身的重放能力它管不了。而让源具备重放能力这件事,在工程上通常只有一个答案——在 Logstash 前面放一个 Kafka(或同类的持久化消息系统)。这是可靠性设计的第一个岔路口:把缓冲和重放放在 Logstash 进程内,还是放在它外面。
1 | |
两者的取舍不在性能而在边界位置:
| 维度 | 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 | |
这套模式和 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 也会在判定停滞后主动强制退出,这几种情况下内存队列同样全丢。受控重启的安全性取决于关闭流程能不能走完,不取决于它名义上是不是"正常退出"。
练习
-
按照本文实验步骤,在本地用
kill -9模拟崩溃,验证 PQ 的重放行为。记录哪些 event 被重放(未 ACK 的),哪些没有重放(已 ACK 的)。把观察结果和 checkpoint 文件的修改时间对照。 -
用内存队列做同样的实验(
queue.type: memory,同样kill -9),确认崩溃时的 event 丢失行为。然后改成用 Ctrl-C(SIGTERM)退出再跑一次,对比两种退出方式下内存队列的差别,并说明为什么queue.drain在这两次实验里都没有参与。 -
思考题:一个"Beats → Logstash(PQ)→ Elasticsearch"的管道,Logstash 崩溃后重启。分析以下三种 event 各自的命运:(a) 已经被 ES output 成功写入并 ACK 的;(b) 在 PQ 里等待被 worker 消费的;© Beats 已发送但 Logstash 崩溃前还没写入 PQ 的。Beats 的 ack 机制在第 3 种情况下如何介入?
系列导航
参考资料
- 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/下的Queue、Checkpoint、Page与io/MmapPageIOV2;logstash-core/lib/logstash/environment.rb是全部 setting 名与默认值的权威来源) - 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 重放产生的重复)
