Logstash 常被当成"把日志灌进 Elasticsearch 的那个工具",学习方式往往是抄一段 logstash.conf、改几个 Grok 正则、跑通就算会了。这条路能用起来,但解释不了几个问题:为什么同样一份配置,pipeline.workers 调大有时提升吞吐、有时毫无变化?为什么下游 Elasticsearch 变慢,会反过来让 input 端读取变慢?持久队列到底在崩溃时保住了什么、又保不住什么?

这些问题的答案都不在插件参数表里,而在两个更底层的抽象上:Logstash 处理的核心对象是 event,处理的骨架是 input-filter-output 三段管道。event 是一个带 @timestamp 的结构化文档,是管道里流动的最小单位;三段管道规定了字节怎么变成 event、event 怎么被加工、加工完的 event 怎么发往下游。codec、Grok、队列、批处理、背压这些机制,全都是围绕"event 在三段管道里怎么流"展开的。

本系列只抓一个问题:Logstash 这套以 event 为核心对象、以三段管道为骨架的流式 ETL 引擎,在 codec、filter、pipeline、队列各个层面是如何具体实现的,每个实现里有哪些可以迁移到其他流处理、消息中间件、ETL 系统的设计模式。

event:Logstash 流动的最小单位

理解 Logstash,从理解它搬运的东西开始,而不是从插件开始。

Logstash 里流动的每一个东西都是一个 event。它不是一行原始文本,而是一个结构化对象,概念上接近一个 JSON 文档,但多了几个 Logstash 约定的字段:

1
2
3
4
5
event
├── @timestamp 事件时间(ISO8601,带时区),几乎所有下游都依赖它
├── @metadata 临时元数据,能参与处理,但默认不进入最终输出
├── message 原始文本(如果来自纯文本输入)
└── <任意字段> filter 阶段切分 / 加工出来的结构化字段

@timestamp 是 event 的一等公民。它不是"日志被处理的时间",默认是 event 进入 Logstash 的时间,但通常会被 date filter 改写成日志里那条记录真正发生的时间——这个改写是很多时序分析正确与否的分界点,第 07 篇会展开。

@metadata 是一块特殊区域。写进 @metadata 的字段能在管道里被后续 filter 和 output 读取、参与条件判断和路由(比如决定写到哪个 ES 索引),但默认不会出现在发往下游的最终文档里。它相当于管道内部的草稿纸。

字段引用有专门的语法。顶层字段写成 [field] 或直接 field,嵌套字段写成 [parent][child]。这套语法在 filter 的条件判断、mutate 操作、output 的动态索引名里到处出现,是读懂任何一份 logstash.conf 的前提。

把 event 当成核心对象,而不是把"日志"当成核心对象,是理解 Logstash 的第一步。一行日志进来,经过 codec 变成 event,经过 filter 长出结构化字段,经过 output 变回下游需要的格式——整条链路都在对 event 做变换。

三段管道:input、filter、output

Logstash 的骨架是三个阶段,顺序固定:

1
2
3
4
       ┌─────────┐     ┌──────────┐     ┌─────────┐     ┌──────────┐
bytes │ input │ ──▶ │ queue │ ──▶ │ filter │ ──▶ │ output │ ──▶ bytes
│ + codec │ │(mem / PQ)│ │ workers │ │ + codec │
└─────────┘ └──────────┘ └──────────┘ └──────────┘

input 负责把数据搬进来:从文件读、监听端口、从 Kafka 消费、接收 Beats 上报。字节流在这一步经过 input 的 codec 转成 event。

进入管道后,event 先落在一个队列里。这个队列默认在内存,也可以配置成落盘的持久队列——这个选择是整个系列可靠性讨论的核心,下面还会提到。

filter 负责加工:Grok 把非结构化文本切成字段、dissect 按分隔符切分、mutate 改字段、date 改写 @timestamp、geoip 补地理信息。filter 由多个 pipeline worker 并发执行,每个 worker 一次处理一批(batch)event。

output 负责发出去:写 Elasticsearch、写文件、发 Kafka。发往下游前经过 output 的 codec 把 event 转回字节。

codec 在这里是一个容易被忽略但很关键的角色。它不属于 filter,而是贴在 input 和 output 边界上的字节↔event 转换器。json codec 把一行 JSON 直接解析成 event 字段,plain 把整行塞进 messagemultiline 把跨行的堆栈合并成一个 event。为什么 multiline 是 codec 不是 filter,第 03 篇会专门讲——简短的答案是:多行合并必须发生在"字节还没被切成一个个 event"的边界上,等切成 event 再想合并就晚了。

队列与背压:可靠性和吞吐都在这里

三段管道之间不是直接函数调用,中间隔着队列,这带来两个后果。

第一个后果是背压。如果 output 端的 Elasticsearch 变慢,output 消费 event 的速度下降,队列积压,filter worker 拿不到空间放新结果也会停下,反过来让 input 减速甚至停止读取。压力就这样从下游一路传回上游。这解释了开头那个现象:下游变慢会让 input 端读取变慢。第 09 篇会把这条背压链路拆开讲。

第二个后果是可靠性的分界线在队列上。默认的内存队列在进程崩溃、断电时,队列里还没被 output 确认发出的 event 会丢失。持久队列(Persistent Queue)把队列落盘并做检查点,在崩溃后重启时可以从盘上恢复未处理完的 event,把可靠性从"尽力而为"提升到 at-least-once。

这里必须说清 at-least-once 的边界,不能简单说成"Logstash 保证不丢数据"。持久队列保证的是"进入队列之后"的 event 在崩溃时不丢;它管不了"还没进队列"的那一段。如果 input 从一个不能重放的数据源读取(比如一个只发一次、不等确认的 UDP 流),数据在进队列之前就可能丢,持久队列无能为力。反过来,能重放的源(Kafka 保留位点、Beats 支持 ack 后重发)配合持久队列,才谈得上端到端的 at-least-once。第 10、11 篇会把队列、检查点、死信队列(DLQ)逐一展开。

关于队列默认值有一个版本敏感的提醒。队列类型由 queue.type 决定,历史上不同版本的默认值和推荐做法有过调整,本系列在给出具体默认值时会标注适用的 Logstash 版本,不写"Logstash 一直默认内存队列"这类全称断言。可以稳定依赖的判断是:内存队列快但崩溃丢数据,持久队列慢一点但能在崩溃边界上兜底,选哪个取决于你对丢数据的容忍度。

一条主线压住整个系列

把这些观察压成一句话:

1
2
3
4
5
Logstash 是一个以 event 为核心对象、以 input-filter-output
三段管道为骨架的流式 ETL 引擎;codec 决定字节与 event 的
互转,filter(尤其 Grok)把非结构化文本切成字段,pipeline
worker 用批处理并发消费,持久队列在崩溃边界上把可靠性从
"尽力而为"提升到 at-least-once。

后续每篇都在展开这条主线下的一个子问题。核心抽象篇(00-03)讲清 event 模型、JRuby on JVM 的运行形态和 codec 边界。插件三段篇(04-08)展开 input 接入、Grok 与 dissect 的取舍、常用 filter 组合和 Elasticsearch output 的批量写入。管道可靠性篇(09-12)展开 worker/batch/背压、内存队列 vs 持久队列、死信队列和 Multiple Pipelines。运维调优篇(13-14)讲监控指标和性能调优。演进对比篇(15-17)把 Logstash 和 Beats、Ingest Pipeline、Fluentd、Vector、Elastic Agent 放在一起看分工与取舍。

实验:用 rubydebug 看清一个 event

Logstash 的核心抽象在一条最小命令里就能看见。它不需要配置文件、不写任何下游,只把 event 打回终端:

1
bin/logstash -e 'input { stdin {} } filter { } output { stdout { codec => rubydebug } }'

启动后在终端输入一行普通文本,比如 hello logstash,回车。Logstash 会打印出这一行被包装成的 event:

1
2
3
4
5
6
{
"message" => "hello logstash",
"@timestamp" => 2026-08-03T03:00:00.123Z,
"@version" => "1",
"host" => { "hostname" => "your-host" }
}

几件事直接对应上面的抽象。输入的纯文本被 plain codec(stdin 的默认 codec)原样放进了 message 字段,而不是被解析。@timestamp 被自动加上,值是 event 进入 Logstash 的时刻——因为这条管道的 filter 是空的,没有 date filter 去把它改写成"日志本身的时间"。stdoutrubydebug codec 把 event 对象格式化成了这段可读的结构,这就是 output 端 codec 在做的字节转换的一个特例(转成人可读的调试输出)。

再做一个对照实验,把 input 的 codec 换成 json

1
bin/logstash -e 'input { stdin { codec => json } } output { stdout { codec => rubydebug } }'

这次输入一行 JSON,比如 {"level":"ERROR","svc":"pay"},回车。输出里不再有一个装着整行文本的 message 字段,取而代之的是 levelsvc 两个顶层字段。同一段字节,换一个 codec,就决定了它变成一个"只有 message 的 event"还是"带结构化字段的 event"。这一步就是 codec 作为字节↔event 边界转换器的直接体现,也说明了为什么 codec 的选择要放在 input/output 边界上,而不是丢给 filter。

模式提炼

Logstash 的整套设计可以提炼成一个可迁移到其他流处理系统的模式:

1
2
3
4
5
6
7
模式:统一事件对象 + 三段管道 + 队列解耦 + 批处理消费

- 把所有输入归一成同一种事件对象,后续处理只面对这一种结构
- 把"读入 / 加工 / 发出"切成边界清晰的三段,各段只管自己那一步
- 在段与段之间放队列,用队列解耦生产和消费速度,让背压自然传导
- 消费端按批处理并发拉取,用批量摊薄单条处理和网络往返的开销
- 把可靠性的锚点放在队列的持久化和确认上,而不是每段各自保证

这个模式不是 Logstash 独有。Kafka 的消费者按批 poll、用 offset 提交表达"处理到哪了",是同一套思路在消息系统的具体化。Flink 把计算表达成算子(operator)链、用 checkpoint 兜底状态,是同一套思路在流计算引擎的具体化。任何一个"读源 → 转换 → 写目标"的通用 ETL 作业,都能对上 input-filter-output 这三段。理解这套模式,比记住某个 filter 的参数更值得带走。

工程迁移表

Logstash 概念 Kafka 生态对应 Flink 对应 通用 ETL 对应
event(统一事件对象) ConsumerRecord StreamRecord 一行记录 / 一条消息
input(读入段) Consumer Source Extract
filter(加工段) 应用层处理逻辑 Transformation 算子 Transform
output(发出段) Producer Sink Load
pipeline worker + batch 消费者线程 + poll 批量 算子并行度 + 缓冲 批处理并发
内存 / 持久队列 内存缓冲 / broker 落盘 内存 / RocksDB state 内存 / 落盘 staging
持久队列 at-least-once offset 提交语义 checkpoint + 重放 幂等写 + 重跑
DLQ(死信队列) dead letter topic side output 错误记录表

注意 at-least-once 这一行在各个系统里都带着相同的边界条件:它保证的是"进入受控存储之后"的不丢,端到端不丢还要求源可重放、目标可去重。Logstash 的持久队列、Kafka 的 offset、Flink 的 checkpoint,都不能单独承诺端到端 exactly-once,这一点在第 10 篇会讲透。

常见误解

误解一:“Logstash 就是把日志发进 ES 的工具”。发进 ES 只是 output 的一种。Logstash 的骨架是通用的三段 ETL 管道,output 可以是文件、Kafka、S3、另一个 Logstash 管道。把它等同于"ES 的采集器",会忽略它作为通用流式转换引擎的定位,也解释不了它和 Beats、Fluentd 的竞争关系(第 15、16 篇)。

误解二:“Logstash 保证数据不丢”。这句话不加限定就是错的。默认的内存队列在崩溃时会丢队列里的 event;持久队列只保证"进入队列之后"的 at-least-once,且要求源可重放才谈得上端到端。任何关于 Logstash 可靠性的断言都必须带上队列类型和故障边界。

误解三:“Grok 是一门配置语言,很难学”。Grok 不是新语言,它是命名正则加一套预打包的 pattern 别名。%{IP:client} 展开后就是一段带命名捕获组的正则。写不出匹配时,问题几乎总是正则本身,而不是 Grok。第 05 篇会把这层"别名"揭开。

误解四:“调大 pipeline.workers 一定能提速”。worker 数决定 filter/output 阶段的并发度,但如果瓶颈在 output 下游(比如 ES 写入已经饱和)或在单条 event 的 I/O 等待上,加 worker 只会加剧背压而不提升吞吐。吞吐问题要先定位卡在 input、filter 还是 output,这是第 09、13 篇的主题。

练习

  1. 本地用 Docker 起一个 Logstash(或用已有实例),运行本文第一个 rubydebug 实验,观察一行文本被包装成的 event,找出 message@timestamp@version 三个字段。再把 input codec 换成 json,输入一行 JSON,对比 event 结构的变化。把这个变化和"codec 是字节↔event 边界转换器"这句话对上。

  2. 在上面的管道里加一个 filter:filter { mutate { add_field => { "[@metadata][tag]" => "test" } add_field => { "real_field" => "x" } } }。观察 rubydebug 输出里 real_field 出现了、而 @metadata 里的字段没有出现在最终输出中。解释为什么 @metadata 能被处理却默认不输出。

  3. 思考题:给定一个"从 UDP 端口收 syslog、Grok 解析后写入 ES"的管道,分别指出在哪些环节可能丢数据(提示:UDP 本身、进队列前、内存队列崩溃、ES 写入失败)。哪些环节持久队列能兜底,哪些不能?把答案和第 10 篇的 at-least-once 边界对照。

系列导航

序号 主题 状态
00 导读:核心对象是 event,骨架是三段管道 本篇
01 Logstash 架构:JRuby、JVM 与 pipeline 的运行形态 下一篇
02 event 模型:@timestamp@metadata 与字段引用
03 codec:字节流与 event 的边界转换
04-08 插件三段(input / Grok / dissect / mutate·date·geoip / ES output) 后续阶段
09-12 管道执行与可靠性(worker·batch·背压 / 队列 / DLQ / Multiple Pipelines) 后续阶段
13-14 运维、监控与调优(Node Stats 与瓶颈定位 / 性能调优) 后续阶段
15-17 演进、生态与对比(vs Beats·Ingest / vs Fluentd·Vector / Elastic Agent 冲击) 后续阶段

Logstash 的 output 最常见的下游是 Elasticsearch,涉及 bulk 写入、索引模板、数据流的地方会交叉引用仓库内《深入 Elasticsearch》系列。

参考资料