深入 Logstash 11 - 死信队列(DLQ):无法处理的 event 去哪
上一篇讲清了持久队列如何在崩溃边界上把 at-least-once 的保证落盘。这一篇进入另一类失败:event 本身在逻辑上就无法被 output 接受,不管重试多少次都不会成功。这类 event 的归宿是死信队列(Dead Letter Queue,DLQ)。
核心问题是三个:什么样的失败会写进 DLQ,DLQ 里的 entry 长什么样,以及怎么读出来重新处理。
两类失败:重试能解的和重试解不了的
持久队列解决的是"崩溃之后不丢 event",但"不丢"和"能处理"是两回事。一个 event 从 PQ 里出来、送往 Elasticsearch,ES 有可能返回两类错误:
1 | |
DLQ 的用途就是第二类:把那些 output 永久拒绝的 event 存到一个单独的文件存储里,而不是直接丢弃,等待人工介入或自动化的补救管道重新处理。
DLQ 数据流
1 | |
DLQ 是文件存储,不是 Kafka topic,也不是独立进程。它由 Logstash 进程自己维护,位置在 path.data(默认 data/)下,子路径按 pipeline id 分隔:data/dead_letter_queue/{pipeline_id}/。
启用 DLQ
在 logstash.yml 里开关 DLQ:
1 | |
启用之后,支持 DLQ 写入的 output(目前主要是 Elasticsearch output)会在遇到 document-level 永久失败时把 event 写进 DLQ,而不是直接抛弃。Logstash 7.x+ 的 Elasticsearch output 对 DLQ 的支持较完整,旧版本行为有差异,生产使用前核对版本文档。
DLQ Entry 的结构
DLQ 里存的不是裸 event。每个 DLQ entry 是原始 event 加上一段元数据:
1 | |
[dead_letter_queue] 是一个约定的 @metadata 子字段路径,在 DLQ 重处理管道里可以直接读到,用来判断失败原因、决定修复策略。
最小实验:触发 mapping 冲突、观察 DLQ 写入
准备一个会触发 mapping 冲突的场景:假设 ES 里 port 字段已被映射为 integer,但发过去的 event 里 port 是字符串 "unknown"。
1 | |
启用 DLQ 后,这条 event 会因 mapping 冲突被 ES 拒绝、写入 data/dead_letter_queue/main/ 目录下的段文件。
读 DLQ 的管道单独定义,通常配置在 pipelines.yml 里作为第二条 pipeline:
1 | |
重处理管道从 DLQ 读出 entry、修复字段、重新写入 ES。这一次因为 port 已经是 integer,写入成功。
对应回内部对象
从实验对应到 Logstash 内部机制:
dead_letter_queue.enable: true 让 Elasticsearch output 插件在调用 ES bulk API 时,对返回的 errors: true 响应里属于 document-level 失败的条目(response status 400/409/类 mapper_parsing_exception)触发 DLQ 写入路径,而不是走 retry 路径。
DLQ 本身是一个 append-only 的分段文件存储(segment file),写入者(主管道的 output)和读取者(DLQ 重处理管道的 input)通过文件偏移量(类似 Kafka offset)解耦:读取者维护自己的消费位点,不影响写入者继续写。
dead_letter_queue.max_bytes 限制 DLQ 文件存储的总大小。当 DLQ 写满时,最新的 Logstash 版本的行为是丢弃最旧的数据为新 entry 腾空间(drop_newer 或 drop_older 依版本而定),因此 DLQ 并不是一个可靠的无限缓存,只是一个有容量上限的观测窗口。
模式提炼
1 | |
工程迁移表
| Logstash 概念 | Kafka 生态对应 | Flink 对应 | 通用 ETL 对应 |
|---|---|---|---|
| DLQ(死信队列) | Dead Letter Topic | Side Output(错误流) | 错误记录表 / 失败文件 |
| DLQ entry(原始 event + 失败元数据) | DLT 消息(含原始 payload + header) | SideOutput 里的 OutputTag 记录 | 错误表行(含 reason 字段) |
| dead_letter_queue input(DLQ 读取) | 单独消费者消费 DLT | 分支处理错误流 | 错误修复作业 |
| plugin_id / reason 元数据 | Kafka header(来源 topic、异常栈) | OutputTag 类型信息 | reason / source 列 |
| max_bytes 容量限制 | retention policy(time / bytes) | State TTL | 表分区 / 归档策略 |
| DLQ commit_offsets | consumer group offset commit | checkpoint 状态推进 | 错误表处理游标 |
Kafka 生态里 Dead Letter Topic 的思路和 Logstash DLQ 几乎同构:把无法处理的消息发往独立 topic,由独立消费者组处理,原消费者不受阻塞。Flink 的 Side Output 在算子层面做同样的分流——正常记录走主流,异常记录走 side output。共同的设计意图是:把错误记录从主链路隔离出去,保证主链路不被问题数据卡住,同时保留可观测的错误记录供后续修复。
常见误解
误解一:“DLQ 和持久队列是同一回事”。持久队列(PQ)解决的是崩溃时不丢还没处理完的 event,关注的是进程生命周期边界;DLQ 解决的是 event 逻辑上无法被 output 接受的情况,关注的是业务语义边界。两者都是文件存储,但目的和触发条件完全不同。
误解二:“启用 DLQ 之后就不会丢数据了”。DLQ 接收的是 output 永久拒绝的 event,它本身也有容量上限(max_bytes)。DLQ 写满后新 entry 仍然可能被丢弃。DLQ 不是无限缓存,它是一个有容量的观测窗口,需要配套监控和定期消费。
误解三:“所有 output 都支持 DLQ”。截至目前,Logstash 的 DLQ 写入主要由 Elasticsearch output 插件实现。其他 output 插件(file、kafka、s3 等)不会自动把永久失败的 event 写进 DLQ,失败仍然走各自的 retry / dead drop 逻辑。
误解四:“DLQ 重处理管道只是把 event 原样重发”。如果原样重发,同样会触发 mapping 冲突再次写入 DLQ,进入死循环。DLQ 重处理管道的核心职责是先修复 event(rename、convert、drop 问题字段),再写回目标,修复逻辑必须依赖 [dead_letter_queue][reason] 字段里的失败原因来选择处理分支。
练习
-
在本地 ES 里给一个索引强制设置 mapping:
port字段为integer。用 Logstash 的generator输入发一条含"port":"unknown"的 event,开启 DLQ 后观察data/dead_letter_queue/目录下是否生成了段文件。再写一条 DLQ 重处理管道,把port修复后重新写入 ES,观察 DLQ 段文件的消费位点是否推进。 -
修改
dead_letter_queue.max_bytes为一个很小的值(如1mb),持续发送会触发 mapping 冲突的 event,观察 DLQ 写满后的行为:是 drop_newer(丢最新进来的)还是 drop_older(丢最旧的)。核对当前 Logstash 版本的文档确认这个行为的版本边界。 -
思考题:如果同一个 Logstash 进程运行两条管道(pipeline A 和 pipeline B),它们各自有一个
data/dead_letter_queue/A/和data/dead_letter_queue/B/目录。DLQ 重处理管道在dead_letter_queueinput 里需要指定pipeline_id,这个设计选择意味着什么?DLQ 隔离设计和 Multiple Pipelines 之间有什么呼应关系?(提示:下一篇就讲 Multiple Pipelines。)
系列导航
| 序号 | 主题 | 状态 |
|---|---|---|
| 00 | 导读:核心对象是 event,骨架是三段管道 | |
| 01 | Logstash 架构:JRuby、JVM 与 pipeline 的运行形态 | |
| 02 | event 模型:@timestamp、@metadata 与字段引用 |
|
| 03 | codec:字节流与 event 的边界转换 | |
| 04 | input 接入:文件、Beats、Kafka 与 TCP | |
| 05 | Grok:命名正则与 pattern 别名 | |
| 06 | dissect:分隔符切分与 Grok 的取舍 | |
| 07 | 常用 filter 组合:mutate、date、geoip、ruby | |
| 08 | Elasticsearch output:bulk 写入、索引模板与数据流 | |
| 09 | pipeline worker、batch 与背压:吞吐调优的三角 | |
| 10 | 内存队列 vs 持久队列:at-least-once 的边界 | |
| 11 | 死信队列(DLQ):无法处理的 event 去哪 | 本篇 |
| 12 | Multiple Pipelines 与 pipeline-to-pipeline | 下一篇 |
| 13 | 监控与诊断:Node Stats API 与瓶颈定位 | |
| 14 | 性能调优:JVM、worker、batch 的系统性方法 | |
| 15 | 演进一:Logstash vs Beats + Ingest Pipeline | |
| 16 | 演进二:Logstash vs Fluentd、Vector | |
| 17 | 演进三:Elastic Agent 与 Logstash 的分工 |
参考资料
- Logstash 死信队列官方文档:https://www.elastic.co/guide/en/logstash/current/dead-letter-queues.html(DLQ 启用、max_bytes、文件结构、drop 行为的版本说明)
- dead_letter_queue input 插件文档:https://www.elastic.co/guide/en/logstash/current/plugins-inputs-dead_letter_queue.html(path、pipeline_id、commit_offsets、clean_consumed 参数)
- Elasticsearch output 与 DLQ 的集成:https://www.elastic.co/guide/en/logstash/current/plugins-outputs-elasticsearch.html(document-level 失败触发 DLQ 的条件)
- Logstash 持久队列文档:https://www.elastic.co/guide/en/logstash/current/persistent-queues.html(PQ 与 DLQ 的区别和配合关系)
- Logstash 源码仓库(DLQ 实现):https://github.com/elastic/logstash(
logstash-core/src/main/java/org/logstash/common/io/DeadLetterQueueWriter.java) - 仓库内《深入 Elasticsearch》系列:bulk 写入拒绝原因与 mapping conflict 的技术细节
