上一篇讲清了持久队列如何在崩溃边界上把 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.dead_letter_queue 指定,默认取 path.data/dead_letter_queue;再往下按 pipeline id 分隔:data/dead_letter_queue/{pipeline_id}/。这个按 id 分隔的布局在后面讲重处理管道时还会用到。

启用 DLQ

logstash.yml 里开关 DLQ:

1
2
3
4
5
6
# logstash.yml
dead_letter_queue.enable: true
dead_letter_queue.max_bytes: 1gb # 默认 1024mb
dead_letter_queue.storage_policy: drop_older # 默认 drop_newer
dead_letter_queue.retain.age: 5d # 超过 5 天的 entry 自动过期清理
path.dead_letter_queue: /var/lib/logstash/dead_letter_queue # 不写则用 path.data/dead_letter_queue

启用之后,支持 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
2
3
4
5
6
7
DLQ entry
├── 原始 event(全部字段,包含 @timestamp 与原有 @metadata)
└── [@metadata][dead_letter_queue] 元数据段
├── plugin_id 写入失败的 output 插件实例 id
├── plugin_type "output"
├── reason ES 返回的拒绝原因(mapping conflict 描述)
└── entry_time event 进入 DLQ 的时间戳

这四个字段挂在 @metadata 下面,完整路径是 [@metadata][dead_letter_queue][reason] 这种形态。在 DLQ 重处理管道里可以直接读到,用来判断失败原因、决定修复策略。

前缀不能省。@metadata 是 event 的独立命名空间,写成 [dead_letter_queue][reason] 引用的是一个普通顶层字段,而那个字段不存在,条件判断静默为假,分支永远不命中,Logstash 也不会报任何错。这是 DLQ 重处理最容易踩的一个坑,因为它不以失败的形式表现出来,而是表现成"修复逻辑好像没生效"。

另一面是 @metadata 不会被 output 写出。如果要把 reason 存进 ES 供后续排查,得先用 mutate 把它复制到一个普通字段:

1
2
3
4
5
filter {
mutate {
copy => { "[@metadata][dead_letter_queue][reason]" => "dlq_reason" }
}
}

最小实验:触发 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 writer 攒够一批或者等到 dead_letter_queue.flush_interval(默认 5000 毫秒)才把缓冲刷到磁盘,所以 kill 之后马上去看目录,很可能什么都没有。要判断实验是否成功,至少等过一个 flush 周期再列目录。

读 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
23
24
25
26
27
28
29
30
31
# dlq_replay.conf(DLQ 读取与重处理管道,id: "dlq_replay")
input {
dead_letter_queue {
path => "/path/to/data/dead_letter_queue"
pipeline_id => "main"
commit_offsets => true # 默认 true,处理成功后推进读取位置
clean_consumed => true # 默认 false;开启才自动删除已消费的段文件
# 要求 Logstash >= 8.4.0,且必须与 commit_offsets => true 同时使用
}
}
filter {
# 对于 port 为 "unknown" 的情况,把字段转成 -1
if [port] == "unknown" {
mutate { replace => { "port" => -1 } }
mutate { convert => { "port" => "integer" } }
} else {
# 无法自动修复的形态("N/A"、超长字符串、结构完全不同的 event)
# 归档到本地文件人工处理,并打 tag 供告警统计,不再重发
mutate { add_tag => ["dlq_unrepairable"] }
}
}
output {
if "dlq_unrepairable" in [tags] {
file { path => "/var/log/logstash/dlq-unrepairable-%{+YYYY.MM.dd}.log" }
} else {
elasticsearch {
hosts => ["http://localhost:9200"]
index => "app-logs"
}
}
}

重处理管道从 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
curl -s http://localhost:9600/_node/stats/pipelines/main | python3 -m json.tool
1
2
3
4
5
6
pipelines.main.dead_letter_queue.queue_size_in_bytes      # DLQ 当前占用字节数
pipelines.main.dead_letter_queue.max_queue_size_in_bytes # 即 dead_letter_queue.max_bytes
pipelines.main.dead_letter_queue.storage_policy # drop_newer 或 drop_older
pipelines.main.dead_letter_queue.dropped_events # 因写满被丢弃的 entry 数
pipelines.main.dead_letter_queue.expired_events # 因 retain.age 过期被清理的 entry 数
pipelines.main.dead_letter_queue.last_error # 最近一次 DLQ 写入错误

这一组字段只在 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
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,默认 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_consumeddead_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。

练习

  1. 在本地 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 是否被整体重放一遍。

  2. dead_letter_queue.max_bytes 改成一个很小的值(如 1mb),持续发送会触发 mapping 冲突的 event 直到写满。在 dead_letter_queue.storage_policydrop_newerdrop_older 两种配置下各跑一次,对比写满之后新旧 entry 的去留差异,并用 pipelines.main.dead_letter_queue.dropped_events 核对被丢弃的条数。

  3. 在重处理管道里先不开 clean_consumed 跑完一轮,看段文件是否还留在目录里;再打开它重跑一轮(同时保留 commit_offsets => true),对比目录变化。最后故意把 commit_offsets 改成 false 再启动一次,记录报错信息,说明这个约束为什么是硬性的。

  4. 思考题:如果同一个 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 架构: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 的冲击

参考资料