深入 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.dead_letter_queue 指定,默认取 path.data/dead_letter_queue;再往下按 pipeline id 分隔:data/dead_letter_queue/{pipeline_id}/。这个按 id 分隔的布局在后面讲重处理管道时还会用到。
启用 DLQ
在 logstash.yml 里开关 DLQ:
1 | |
启用之后,支持 DLQ 写入的 output(目前主要是 Elasticsearch output)会在遇到 document-level 永久失败时把 event 写进 DLQ,而不是直接抛弃。
上面这份配置里,容量治理是两条腿。storage_policy 管的是"写满了丢谁",属于被动兜底;dead_letter_queue.retain.age(8.4 起可用,写成 5d 这种时长)管的是"多久之前的自动清掉",属于主动回收。只配前者,DLQ 会一直贴着 max_bytes 满着,里面混着几个月前的陈旧 entry;配上后者,过期段文件会被自动删除,磁盘和排查窗口都受控。两条腿都不配,DLQ 目录的增长只能靠人去消费,这是最常见的一种失控形态。
DLQ Entry 的结构
DLQ 里存的不是裸 event。每个 DLQ entry 是原始 event 加上一段元数据:
1 | |
这四个字段挂在 @metadata 下面,完整路径是 [@metadata][dead_letter_queue][reason] 这种形态。在 DLQ 重处理管道里可以直接读到,用来判断失败原因、决定修复策略。
前缀不能省。@metadata 是 event 的独立命名空间,写成 [dead_letter_queue][reason] 引用的是一个普通顶层字段,而那个字段不存在,条件判断静默为假,分支永远不命中,Logstash 也不会报任何错。这是 DLQ 重处理最容易踩的一个坑,因为它不以失败的形式表现出来,而是表现成"修复逻辑好像没生效"。
另一面是 @metadata 不会被 output 写出。如果要把 reason 存进 ES 供后续排查,得先用 mutate 把它复制到一个普通字段:
1 | |
最小实验:触发 mapping 冲突、观察 DLQ 写入
准备一个会触发 mapping 冲突的场景:假设 ES 里 port 字段已被映射为 integer,但发过去的 event 里 port 是字符串 "unknown"。
1 | |
启用 DLQ 后,这条 event 会因 mapping 冲突被 ES 拒绝、写入 data/dead_letter_queue/main/ 目录下的段文件。
段文件不是立刻出现的。DLQ writer 攒够一批或者等到 dead_letter_queue.flush_interval(默认 5000 毫秒)才把缓冲刷到磁盘,所以 kill 之后马上去看目录,很可能什么都没有。要判断实验是否成功,至少等过一个 flush 周期再列目录。
读 DLQ 的管道单独定义,通常配置在 pipelines.yml 里作为第二条 pipeline:
1 | |
重处理管道从 DLQ 读出 entry、修复字段、重新写入 ES。这一次因为 port 已经是 integer,写入成功。
else 那条分支不是凑数的。上面的修复逻辑只覆盖 port == "unknown" 这一种形态,如果非法值换成 "N/A"、换成一个超长字符串、或者压根是另一种结构的 event,if 不命中,event 原样重发,再被 ES 以同样的理由拒绝,然后进入 dlq_replay 这条管道自己的 DLQ 目录。有兜底分支才能把这类 event 拦在归档文件里,同时靠 tag 计数暴露出来。
对应回内部对象
从实验对应到 Logstash 内部机制:
dead_letter_queue.enable: true 让 Elasticsearch output 插件在调用 ES bulk API 时,对返回的 errors: true 响应逐条按状态码分流,而不是一律走 retry 路径。分流规则是:
- 200 / 201:写入成功。
- 400 / 404:document-level 永久失败(mapping 冲突、目标索引不存在等),写进 DLQ。插件源码里这两个码就是一个常量:
DOC_DLQ_CODES = [400, 404]。 - 409(version conflict):记一条日志后直接丢弃,既不重试也不进 DLQ。
- 其余状态码(429、503 等):走 retry 路径。
400/404 这个集合可以用 ES output 的 dlq_custom_codes 参数扩展,把别的状态码也纳入 DLQ 路由。但它不能覆盖成功码,也不能包含 409:配置里写了 409,插件在启动阶段就抛 ConfigurationError。
DLQ 本身是一个 append-only 的分段文件存储(segment file),写入者(主管道的 output)和读取者(DLQ 重处理管道的 input)通过文件偏移量(类似 Kafka offset)解耦:读取者维护自己的消费位点,不影响写入者继续写。这个位点存在 sincedb 文件里,默认落在 path.data/plugins/inputs/dead_letter_queue 下,文件名由 DLQ 路径的哈希决定。删掉它就等于把消费位置重置到段文件的开头,整个 DLQ 会被重放一遍。
dead_letter_queue.max_bytes 限制 DLQ 文件存储的总大小,默认 1024mb。写满之后丢谁,由 dead_letter_queue.storage_policy 决定,取值只有两个:
drop_newer(默认):丢新来的 entry,保住已经存下的。drop_older:丢最旧的 entry,为新 entry 腾空间。
默认落在 drop_newer 这一侧,意味着 DLQ 一旦写满,后续失败的 event 会直接消失,而不是把旧记录挤出去。这个方向对排查的影响不小:最早出现的那批故障被完整保留,之后的看不到。想要"总能看到最近的失败",必须显式改成 drop_older。不管取哪个,DLQ 都不是可靠的无限缓存,只是一个有容量上限的观测窗口。
监控 DLQ
DLQ 在 Node Stats 里有一等命名空间,per-pipeline 挂在 pipelines.<id>.dead_letter_queue 下:
1 | |
1 | |
这一组字段只在 dead_letter_queue.enable: true 时注册。三条落地的判读规则:
dropped_events > 0已经是事故而非预警:DLQ 写满了,失败记录正在丢。这个计数器数的正是"DLQ 写满之后新 entry 去哪了"的答案。queue_size_in_bytes / max_queue_size_in_bytes当水位用,超过某个比例(比如 70%)就该确认有没有人在消费这个 DLQ。expired_events持续增长说明retain.age在生效,同时也说明失败 entry 从来没被处理过就过期了,说明修复流程在空转。
storage_policy 那一格附带一个用处:不用翻配置文件就能确认当前实例实际跑的是 drop_newer 还是 drop_older。
模式提炼
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,默认 1024mb)。默认 drop_newer 之下,写满后新来的失败 event 直接消失。dropped_events 这个指标就是为这件事存在的,它大于 0 意味着已经在丢了。DLQ 不是无限缓存,是一个有容量的观测窗口。
误解三:“所有 output 都支持 DLQ”。截至目前,Logstash 的 DLQ 写入主要由 Elasticsearch output 插件实现。其他 output 插件(file、kafka、s3 等)不会自动把永久失败的 event 写进 DLQ,失败仍然走各自的 retry / dead drop 逻辑。
误解四:“已经消费过的 DLQ 段文件会自动清理”。dead_letter_queue input 的 clean_consumed 默认是 false:位点推进了,段文件还留在磁盘上。要自动删除必须显式开启,而且要求 Logstash 8.4.0 以上、必须与 commit_offsets => true 同时使用,否则启动直接报 ConfigurationError。默认值这一侧意味着 DLQ 目录会持续增长,光有重处理管道不等于磁盘受控,真正管回收的是 clean_consumed 和 dead_letter_queue.retain.age 这两个开关。
误解五:“DLQ 重处理管道只是把 event 原样重发”。原样重发会再次触发同一个 mapping 冲突,而失败的 entry 并不会回到原来那个 DLQ——重处理管道有自己的 pipeline id,它写出的失败 entry 落进 dead_letter_queue/<replay_id>/,而它读的是 pipeline_id => "main"。所以后果不是字面意义的死循环,而是更隐蔽的一种:这些 event 掉进了第二个 DLQ 目录,而那个目录没有任何管道在读,只会静静堆积。重处理管道的核心职责是先按 [@metadata][dead_letter_queue][reason] 分支修复 event(rename、convert、drop 问题字段),再写回目标;同时应当显式规划自己的 DLQ 归属,或者干脆用 else 分支把修不了的 event 归档掉,不让它们进入第二层 DLQ。
练习
-
在本地 ES 里给一个索引强制设置 mapping:
port字段为integer。用 Logstash 的generator输入发一条含"port":"unknown"的 event,开启 DLQ 后等过一个dead_letter_queue.flush_interval(默认 5 秒)再列data/dead_letter_queue/main/目录,确认段文件已经生成。再写一条 DLQ 重处理管道把port修好写回 ES,去path.data/plugins/inputs/dead_letter_queue下找到 sincedb 文件,确认消费位点已经推进;然后把这个文件删掉重启重处理管道,观察 DLQ 是否被整体重放一遍。 -
把
dead_letter_queue.max_bytes改成一个很小的值(如1mb),持续发送会触发 mapping 冲突的 event 直到写满。在dead_letter_queue.storage_policy取drop_newer和drop_older两种配置下各跑一次,对比写满之后新旧 entry 的去留差异,并用pipelines.main.dead_letter_queue.dropped_events核对被丢弃的条数。 -
在重处理管道里先不开
clean_consumed跑完一轮,看段文件是否还留在目录里;再打开它重跑一轮(同时保留commit_offsets => true),对比目录变化。最后故意把commit_offsets改成false再启动一次,记录报错信息,说明这个约束为什么是硬性的。 -
思考题:如果同一个 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。)
系列导航
参考资料
- Logstash 死信队列官方文档:https://www.elastic.co/guide/en/logstash/current/dead-letter-queues.html(DLQ 启用、
max_bytes、storage_policy、retain.age、文件结构) - 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、sincedb_path参数与默认值) - Elasticsearch output 与 DLQ 的集成:https://www.elastic.co/guide/en/logstash/current/plugins-outputs-elasticsearch.html(document-level 失败触发 DLQ 的条件、
dlq_custom_codes) - 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 的技术细节
