上一篇讲清了持久队列如何在崩溃边界上把 at-least-once 的保证落盘。这一篇进入另一类失败:event 本身在逻辑上就无法被 output 接受,不管重试多少次都不会成功。这类 event 的归宿是死信队列(Dead Letter Queue,DLQ)。

核心问题是三个:什么样的失败会写进 DLQ,DLQ 里的 entry 长什么样,以及怎么读出来重新处理。

两类失败:重试能解的和重试解不了的

持久队列解决的是"崩溃之后不丢 event",但"不丢"和"能处理"是两回事。一个 event 从 PQ 里出来、送往 Elasticsearch,ES 有可能返回两类错误:

1
2
3
4
5
6
7
8
瞬时失败(transient):
ES 集群暂时不可达、写入超时、节点临时过载
→ 重试就有机会成功,Logstash 的 retry 逻辑能处理

永久失败(permanent / document-level):
mapping 冲突(往 integer 字段写了 string 值)
index 的 mapping 不接受这种文档结构
→ 不管重试多少次,问题出在 event 本身,重试无意义

DLQ 的用途就是第二类:把那些 output 永久拒绝的 event 存到一个单独的文件存储里,而不是直接丢弃,等待人工介入或自动化的补救管道重新处理。

DLQ 数据流

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
input → filter → output(Elasticsearch)

mapping 冲突 / document rejected


┌──────────────────────┐
│ Dead Letter Queue │
│ (file-based store) │
│ data/dlq/{pipeline} │
└──────────────────────┘

(另一条 pipeline 读取)


dead_letter_queue input

fix / transform


output → Elasticsearch

DLQ 是文件存储,不是 Kafka topic,也不是独立进程。它由 Logstash 进程自己维护,位置在 path.data(默认 data/)下,子路径按 pipeline id 分隔:data/dead_letter_queue/{pipeline_id}/

启用 DLQ

logstash.yml 里开关 DLQ:

1
2
3
4
# logstash.yml
dead_letter_queue.enable: true
dead_letter_queue.max_bytes: 1gb
# path.data 指向 data/ 目录,DLQ 文件在其下的 dead_letter_queue/ 子目录

启用之后,支持 DLQ 写入的 output(目前主要是 Elasticsearch output)会在遇到 document-level 永久失败时把 event 写进 DLQ,而不是直接抛弃。Logstash 7.x+ 的 Elasticsearch output 对 DLQ 的支持较完整,旧版本行为有差异,生产使用前核对版本文档。

DLQ Entry 的结构

DLQ 里存的不是裸 event。每个 DLQ entry 是原始 event 加上一段元数据:

1
2
3
4
5
6
7
DLQ entry
├── 原始 event(全部字段,包含 @timestamp、@metadata)
└── [dead_letter_queue] 元数据段
├── plugin_id 写入失败的 output 插件实例 id
├── plugin_type "output"
├── reason ES 返回的拒绝原因(mapping conflict 描述)
└── entry_time event 进入 DLQ 的时间戳

[dead_letter_queue] 是一个约定的 @metadata 子字段路径,在 DLQ 重处理管道里可以直接读到,用来判断失败原因、决定修复策略。

最小实验:触发 mapping 冲突、观察 DLQ 写入

准备一个会触发 mapping 冲突的场景:假设 ES 里 port 字段已被映射为 integer,但发过去的 event 里 port 是字符串 "unknown"

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
# logstash.conf(主处理管道,id: "main")
input {
generator {
lines => ['{"host":"web-01","port":"unknown","status":200}']
count => 1
codec => json
}
}
filter {}
output {
elasticsearch {
hosts => ["http://localhost:9200"]
index => "app-logs"
}
}

启用 DLQ 后,这条 event 会因 mapping 冲突被 ES 拒绝、写入 data/dead_letter_queue/main/ 目录下的段文件。

读 DLQ 的管道单独定义,通常配置在 pipelines.yml 里作为第二条 pipeline:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
# dlq_replay.conf(DLQ 读取与重处理管道)
input {
dead_letter_queue {
path => "/path/to/data/dead_letter_queue"
pipeline_id => "main"
# commit_offsets: true 默认,处理成功后推进读取位置
# clean_consumed: true 处理完的段文件自动清理
}
}
filter {
# 对于 port 为 "unknown" 的情况,把字段转成 -1
if [port] == "unknown" {
mutate { replace => { "port" => -1 } }
mutate { convert => { "port" => "integer" } }
}
}
output {
elasticsearch {
hosts => ["http://localhost:9200"]
index => "app-logs"
}
}

重处理管道从 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
2
3
4
5
6
7
8
9
10
11
模式:死信队列 = 永久失败的隔离存储 + 可观测的失败原因 + 可重处理的补救通道

核心区分:
retry(重试) → 失败原因在外部系统(临时不可达),event 本身没问题
DLQ(死信) → 失败原因在 event 本身(结构、类型与目标 schema 冲突)

设计原则:
- DLQ 不是 retry 的替代,是 retry 穷尽后的兜底
- DLQ entry 必须携带失败元数据,否则重处理时无法判断修复策略
- DLQ 读取管道与主管道相互独立,各自有自己的资源限额和生命周期
- DLQ 有容量上限,监控 DLQ 写入速率是运维的必要指标

工程迁移表

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] 字段里的失败原因来选择处理分支。

练习

  1. 在本地 ES 里给一个索引强制设置 mapping:port 字段为 integer。用 Logstash 的 generator 输入发一条含 "port":"unknown" 的 event,开启 DLQ 后观察 data/dead_letter_queue/ 目录下是否生成了段文件。再写一条 DLQ 重处理管道,把 port 修复后重新写入 ES,观察 DLQ 段文件的消费位点是否推进。

  2. 修改 dead_letter_queue.max_bytes 为一个很小的值(如 1mb),持续发送会触发 mapping 冲突的 event,观察 DLQ 写满后的行为:是 drop_newer(丢最新进来的)还是 drop_older(丢最旧的)。核对当前 Logstash 版本的文档确认这个行为的版本边界。

  3. 思考题:如果同一个 Logstash 进程运行两条管道(pipeline A 和 pipeline B),它们各自有一个 data/dead_letter_queue/A/data/dead_letter_queue/B/ 目录。DLQ 重处理管道在 dead_letter_queue input 里需要指定 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 的分工

参考资料