上一篇讲完了 dissect 如何用分隔符驱动的线性扫描替代 Grok 的回溯正则,代价是无法处理非均匀格式。这一篇进入 filter 阶段最常见的后续加工链:mutate 改写字段、date 把日志时间戳覆盖 @timestamp、geoip 用 IP 换取地理信息,以及用条件判断和 tag 做分支路由。

核心问题:date filter 为什么一定要改写 @timestamp,而不是新增一个字段;条件判断和 tag 驱动的分支是 Logstash 唯一的路由原语,它和消息系统的 topic 路由有何本质区别。

数据流:event 经过 filter 链的变换过程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
           ┌──────────────────────────────────────────────────────┐
│ filter block │
│ │
event in ──▶ grok/dissect ──▶ mutate ──▶ date ──▶ geoip ──▶ event out
│ │ │ │
│ (parse text (overwrite │
│ into fields) @timestamp) │
│ │
│ if [type] == "access" { │
│ geoip { ... } │
│ } else if "_grokparsefailure" in [tags] { │
│ mutate { add_tag => ["parse_error"] } │
│ } │
└──────────────────────────────────────────────────────┘

filter block 是一个顺序执行的指令序列。所有 filter 插件和条件块从上到下依次运行,每一步的输出是下一步的输入。没有隐式的并行——并行发生在 pipeline worker 层面,而不在单个 event 的处理链里。

关键对象

filter 阶段操作的对象始终是 event 的字段映射(field map)。每个 filter 插件本质上是一个 event.fields → event.fields 的函数,读取已有字段,写入或删除字段。

四个常用插件各自的职责:

1
2
3
4
mutate  ── 字段的增删改:rename、replace、convert、gsub、add_field、remove_field
date ── 解析字段里的时间字符串,覆盖写入 @timestamp
geoip ── 查 MaxMind 数据库,把 IP 地址展开为地理信息字段
条件块 ── if/else if/else + 字段引用 + 运算符,决定对哪些 event 执行哪些 filter

tag 是 event 里 [tags] 数组的元素。Grok 解析失败时自动打上 _grokparsefailure,date 解析失败时打上 _dateparsefailure。tag 是条件分支最常见的路由信号:不是在配置里写死"哪类日志走哪个分支",而是让前序 filter 把结果写进 tag,后续条件读 tag 决定走哪条路。

mutate:字段的基本外科手术

mutate 的操作分几类,对应不同的 event 修改意图:

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
filter {
mutate {
# 重命名字段
rename => { "host" => "source_host" }

# 替换字段值(支持字段引用)
replace => { "message" => "processed: %{message}" }

# 类型转换:string → integer / float / boolean
convert => { "response_code" => "integer"
"bytes" => "integer" }

# 正则替换字段内容(全局替换)
gsub => [ "url", "\?.*", "" ]

# 增加字段
add_field => { "env" => "production" }

# 删除不再需要的字段,减少 ES 存储
remove_field => [ "message", "rawlog" ]

# 增加 tag
add_tag => [ "processed" ]
}
}

convert 的重要性容易被低估。Grok 切出来的字段全是字符串,写进 ES 之后如果字段映射是 keyword,数值聚合会失败。convert 在 filter 阶段就把字符串强转成数字,是 ES 索引映射正确的前提。

gsub 接受三元组列表:[字段名, 正则, 替换字符串],可以连续列多组。它在字段内容里做全局正则替换,常见用途是从 URL 里去掉查询参数、清洗特殊字符。

remove_field 在写 ES 之前清理掉 message(原始文本)和中间字段,能显著降低单文档大小和 ES 的 _source 存储压力。

date:@timestamp 为何必须被改写

1
2
3
4
5
6
7
filter {
date {
match => [ "log_time", "dd/MMM/yyyy:HH:mm:ss Z" ]
target => "@timestamp"
timezone => "Asia/Shanghai"
}
}

默认情况下,@timestamp 是 event 进入 Logstash 的时刻,不是日志里那条记录真正发生的时刻。两者之间的差距取决于管道的延迟——正常情况下差几秒,批量补录历史日志时可能差几小时甚至几天。

如果不改写 @timestamp,Kibana 的时间线查询(range query on @timestamp)会按处理时间而非事件时间过滤,导致历史日志全部堆在处理时刻附近、事件的真实时序完全丢失。Elasticsearch 的 rollover/ILM 按 @timestamp 决定文档落在哪个索引分片,用处理时间同样会破坏基于时间的索引分布。

match 数组的第一个元素是存放时间字符串的字段名,后面跟一个或多个格式串。多个格式串是或关系——Logstash 依次尝试,第一个匹配成功就用,全部失败则打 _dateparsefailure tag。格式串使用 Joda-Time 语法(Logstash 运行在 JVM 上,时间解析底层是 Java 库)。

target 默认就是 @timestamp,也可以改写到另一个字段(做时间格式标准化而不替换主时间戳)。timezone 用于源日志不带时区偏移时的补全,不写则用系统时区,这是生产中常见的坑:本地测试正常,换时区的服务器就偏移 8 小时。

date filter 执行完之后,log_time 字段并不会被自动删掉,需要在 mutate 里 remove_field 手动清理,否则 ES 里会同时存原始字符串和标准化后的 @timestamp

geoip:IP → 地理信息字段

1
2
3
4
5
6
7
filter {
geoip {
source => "client_ip"
target => "geoip"
fields => [ "city_name", "country_code2", "latitude", "longitude" ]
}
}

geoip 内部持有一份 MaxMind GeoLite2 数据库(mmdb 格式),在 filter 阶段做内存内查表,不涉及网络请求。source 是存放 IP 的字段,target 是写入地理信息的嵌套字段名(默认 geoip),fields 限制只展开需要的子字段,不写则展开全部(包含 AS 号、邮编等,字段数量多)。

内网 IP 和保留地址段(RFC1918、loopback、link-local)在 MaxMind 库里没有记录,查询会失败,不打 tag 但 target 字段为空。如果日志里的 IP 来自 NAT 后的内网,需要在 geoip 前先用 mutate 或条件块跳过这类地址。

数据库更新频率是运维侧的持续任务。MaxMind 每月发布更新,Logstash 7.x 起支持自动从 MaxMind 账户下载最新 mmdb,配置 database 参数指向本地路径时则完全手动管理。地理信息精度(尤其城市级别)随数据库新旧变化,不要把 geoip 的城市字段当作精确坐标使用。

条件块与 tag 驱动的分支

Logstash filter block 内可以写 if / else if / else 条件块,条件里可以引用 event 字段、检查 tag、使用比较和正则运算符:

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
32
33
filter {
# 先解析
grok {
match => { "message" => "%{COMBINEDAPACHELOG}" }
}

# 按解析结果分支
if "_grokparsefailure" in [tags] {
mutate {
add_tag => [ "raw_unparsed" ]
add_field => { "parse_status" => "failed" }
}
} else {
# 只有解析成功才做后续加工
date {
match => [ "timestamp", "dd/MMM/yyyy:HH:mm:ss Z" ]
target => "@timestamp"
}
mutate {
convert => { "response" => "integer"
"bytes" => "integer" }
remove_field => [ "timestamp" ]
}
if [response] >= 500 {
mutate { add_tag => [ "server_error" ] }
} else if [response] >= 400 {
mutate { add_tag => [ "client_error" ] }
}
if [clientip] =~ /^(?!10\.|192\.168\.|172\.(1[6-9]|2\d|3[01])\.)/ {
geoip { source => "clientip" }
}
}
}

条件操作符覆盖:

1
2
3
4
==  !=  <  >  <=  >=          比较(字符串或数字)
=~ !~ 正则匹配 / 不匹配
in not in 成员检查(tag 数组、字段值)
and or not 逻辑运算

字段引用用 [field_name] 语法,嵌套字段用 [parent][child],不存在的字段在条件里求值为 nil(falsy),不会报错。

tag 驱动分支的核心价值是让前序步骤的结果(成功/失败、类型标识)成为后续步骤的路由信号,而不需要在配置里硬编码"哪类日志走哪条路"。_grokparsefailure_dateparsefailure 是 Logstash 约定的内置信号,其他 tag 可以由 mutate 自由添加,形成应用层自定义的路由语义。

filter 里的条件块不能替代 Multiple Pipelines——条件只在同一个 pipeline 内分支,不能把 event 发到另一个 pipeline。跨 pipeline 路由要用 pipeline input/output(第 12 篇)。

实验:一条 nginx access log 经过完整 filter 链

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
32
33
34
# logstash.conf
input {
stdin {}
}

filter {
grok {
match => {
"message" => '%{IPORHOST:clientip} - %{USER:ident} \[%{HTTPDATE:log_ts}\] "%{WORD:verb} %{URIPATHPARAM:request} HTTP/%{NUMBER:httpver}" %{NUMBER:response} %{NUMBER:bytes}'
}
}

if "_grokparsefailure" not in [tags] {
date {
match => [ "log_ts", "dd/MMM/yyyy:HH:mm:ss Z" ]
target => "@timestamp"
timezone => "UTC"
}
mutate {
convert => { "response" => "integer" "bytes" => "integer" }
remove_field => [ "log_ts", "message" ]
}
if [clientip] !~ /^(10\.|192\.168\.|127\.)/ {
geoip { source => "clientip" target => "geo" fields => ["country_code2","city_name"] }
}
if [response] >= 500 {
mutate { add_tag => ["5xx"] }
}
}
}

output {
stdout { codec => rubydebug }
}

输入一行:

1
203.0.113.42 - frank [10/Aug/2025:13:55:36 +0000] "GET /api/v1/users HTTP/1.1" 200 1234

rubydebug 输出中可以观察到:@timestamp 已从进入时刻改写为 2025-08-10T13:55:36.000Zresponsebytes 是整数;log_tsmessage 已被删除;geo.country_code2 存在(203.0.113.42 是 TEST-NET 保留段,实际会查到空或错误,生产用真实 IP)。_grokparsefailure 不在 tags 里,条件分支全部进了 else 路径。

模式提炼

1
2
3
4
5
6
7
8
模式:tag 驱动的渐进式加工链

- 前序 filter 在 tag 里写入结果信号(成功/失败/类型),
后续 filter 读 tag 决定是否执行,形成懒求值的加工链
- date filter 改写 @timestamp 是时序正确性的必要步骤,不是可选项
- mutate.convert 是"切出字段"和"字段有正确类型"之间的必经桥梁
- geoip 是查表操作,不是网络请求,延迟来自内存查找,不来自外部 IO
- 条件块是 pipeline 内路由,不等于 pipeline 间路由

这套"tag 即信号"的模式可以迁移:Kafka 里用 header 传递处理结果标志,Flink 里用 side output 把失败记录分流,都是同一种"前序写信号、后续读信号"的思路。

工程迁移表

Logstash filter 概念 Kafka 生态对应 Flink 对应 通用 ETL 对应
mutate.convert Consumer 侧反序列化类型断言 TypeInformation 类型绑定 类型转换列 / CAST
date filter 改写 @timestamp 消息 header 里的 event time assignTimestampsAndWatermarks ETL 里的 event_time 列
geoip Flink 异步 IO / broadcast state 查表 RichFunction 加载外部字典 维表 JOIN / lookup join
tag 驱动分支 header/key 路由到不同 topic side output / union CASE WHEN / 条件路由
_grokparsefailure → DLQ 解析失败发 dead letter topic side output 收错误记录 错误记录表 / error log
filter 顺序执行 KStream DSL 链式算子 算子拓扑顺序 SQL 执行计划顺序

date filter 那一行尤其值得注意:Flink 里用 assignTimestampsAndWatermarks 从消息体提取事件时间,Logstash 的 date filter 做的是同一件事。两者都在"数据入流"之后、"时序计算"之前,把处理时间替换成事件时间——这是所有流处理系统里时序正确性的共同前提。

常见误解

误解一:“date filter 只是把时间字符串转成另一个格式存起来”。date filter 的核心动作是覆盖写 @timestamp,把 event 的主时间轴从处理时间切换成事件时间。如果只是想保存格式化后的时间字符串,用 mutate 或 ruby filter 就够,不需要 date filter。

误解二:“geoip 会拖慢管道,因为要查外部服务”。geoip 数据库在 Logstash 启动时加载进内存,查询完全在进程内完成,没有网络 IO。性能瓶颈是内存读取和 MaxMind 数据库的查找开销,与外部服务无关。数据库文件越大,冷启动加载时间越长,但单次查询延迟通常在微秒级。

误解三:“条件块里的字段引用如果字段不存在会报错”。Logstash 条件里引用不存在的字段会得到 nil,等价于 false,不抛异常。这意味着 if [nonexistent_field] == "value" 在字段缺失时安静地走 else 分支,不会中断 pipeline。副作用是字段名拼写错误不会有任何提示,需要用 rubydebug 实验确认字段是否存在。

误解四:“filter 里的多个 mutate 操作可以合并成一个”。mutate 内部的多个操作有固定的执行顺序:rename → update → replace → merge → coerce → split → join → gsub → lowercase → uppercase → strip → remove → add_tag → remove_tag → add_field → remove_field。这个顺序与书写顺序无关。如果需要先 rename 再基于新名字做 gsub,把它们写在同一个 mutate 里,rename 总在 gsub 之前执行,符合预期;但如果需要先 add_field 再 rename 同名字段,同一个 mutate 里的顺序反而是 rename 先于 add_field,会踩坑。遇到操作顺序敏感的情况,拆成两个 mutate 块并明确书写顺序。

练习

  1. 用 rubydebug 管道验证 date filter 改写效果:构造一条带历史时间戳的日志(比如一个月前的时间),分别观察有 date filter 和没有 date filter 时 @timestamp 的差异。确认 Kibana 的时间线查询在两种情况下会落在哪个时间段。

  2. 在 filter 里加一个 mutate,对同一字段先做 rename、再做 gsub,然后把它们拆成两个 mutate 块,对比行为是否一致。阅读 mutate 文档里关于操作执行顺序的说明,解释结果。

  3. 构造一条解析失败的日志(让 Grok 匹配不上),观察 tags 数组里出现 _grokparsefailure;然后在 filter 里加条件块,对 _grokparsefailureadd_tag => ["needs_manual_review"],并 add_field => { "parse_status" => "failed" }。验证 rubydebug 输出里两个 tag 都存在。

系列导航

序号 主题 状态
00 导读:核心对象是 event,骨架是三段管道 已发布
01 Logstash 架构: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-14 运维、监控与调优
15-17 演进、生态与对比

参考资料