上一篇讲完了 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] 数组的元素。Logstash 约定了一批内置 tag:Grok 解析失败打 _grokparsefailure,date 解析失败打 _dateparsefailure,geoip 查不到 IP 打 _geoip_lookup_failure,geoip 数据库因长期未更新而停止 enrich 时打 _geoip_expired_database。这些 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
26
27
filter {
mutate {
# 重命名字段
rename => { "host" => "source_host" }

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

# 类型转换:合法目标类型是 integer / integer_eu / float / float_eu
# / string / boolean
convert => { "response_code" => "integer"
"bytes" => "integer" }

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

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

# 注意:remove_field 最后执行,所以上面对 message 的 replace 会被这行抹掉。
# 这正是下文"误解一"要讲的顺序陷阱
remove_field => [ "message", "rawlog" ]

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

这段配置故意保留了一个顺序冲突:replacemessage 里写了 processed: ...remove_field 又把 message 删掉,两个操作还在同一个 mutate 块里。结果是输出里连 message 字段都不存在,那次 replace 完全白做。想保住 replace 的结果,就得把 remove_field 拆到后面一个独立的 mutate 块里。为什么是这个顺序,下文"误解一"里展开。

Grok 切出来的字段默认全是字符串,写进 ES 之后如果映射成了 keyword,数值聚合会失败。把字符串变成数字有两条路径:一条是 grok 的内联类型后缀 %{NUMBER:response:int},在切字段的同时就完成转换;另一条是 mutate 的 convert,在切完之后单独转一遍。内联后缀省一步,但官方写明它只支持 intfloat 两种转换。需要 boolean,或者要解析欧洲数字格式(1.234,56 这种拿逗号当小数点的写法,对应 integer_eu / float_eu),就只有 convert 能做。这是它无法被内联后缀替代的地方。

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 库)。

除了自定义格式串,match 还接受四个内置的格式字面量,覆盖了实际日志里最常见的几种时间形态:

1
2
3
4
ISO8601   ── 任何合法的 ISO8601 时间戳,如 2011-04-19T03:44:01.123456789Z
UNIX ── epoch 秒,整数或浮点都能解析,如 1326149001 / 1326149001.132
UNIX_MS ── epoch 毫秒,整数,如 1366125117000
TAI64N ── tai64n 格式

这四个里前两个在实际日志里出现频率最高。结构化日志框架(Logback 的 JSON encoder、zap、结构化 slog)默认输出 ISO8601,各类 SDK 与消息中间件的时间字段普遍是 epoch 毫秒。遇到这两种直接写字面量,不必手写格式串。

格式串的解析库还受 precision 影响。默认 ms,走 joda-time;设成 ns 之后时间戳保留纳秒精度,解析改由 java.time 承担。两套语法大体相同但有差异,比如时区 ID 在 java.time 里写 VV,在 joda-time 里写 ZZZ。改 precision 时要连带检查格式串。

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

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

geoip:IP → 地理信息字段

1
2
3
4
5
6
7
8
filter {
geoip {
source => "client_ip"
target => "geoip"
fields => [ "city_name", "country_code2", "latitude", "longitude" ]
ecs_compatibility => disabled # 本文所有 geoip 示例都固定在这个模式
}
}

geoip 内部持有 MaxMind GeoLite2 数据库(mmdb 格式),在 filter 阶段做内存内查表,不涉及网络请求。source 是存放 IP 的字段,target 是写入地理信息的嵌套字段名,fields 限制只展开需要的子字段。

插件同时捆绑了 GeoLite2-City 和 GeoLite2-ASN 两个库,default_database_type 默认 City。不写 fields 时展开的是当前选定那个库的全部字段,City 库给的是城市、国家、大洲、邮编、经纬度、时区、region 这一组,字段数量确实不少。但 ASN 号和 AS 组织名不在 City 库里。要拿这两个字段,得把 default_database_type 设成 ASN,或者用 database 指向自己下载的 GeoLite2-ASN.mmdb。这也意味着一个 geoip 实例只查一个库,同时要地理位置和 ASN 就得写两个 geoip 块。

字段落位取决于 ECS 模式

上面那句 ecs_compatibility => disabled 不是可省略的装饰。Logstash 8 里所有插件默认跑 ECS v8 模式(pipeline.ecs_compatibility 的默认值就是 v8),而 geoip 在两种模式下的输出形状完全不同:

数据库字段名 disabled(写在 target 下) ECS v1/v8
city_name [geoip][city_name] [<target>][geo][city_name]
country_code2 [geoip][country_code2] [<target>][geo][country_iso_code]
country_name [geoip][country_name] [<target>][geo][country_name]
postal_code [geoip][postal_code] [<target>][geo][postal_code]
latitude [geoip][latitude] [<target>][geo][location][lat]
longitude [geoip][longitude] [<target>][geo][location][lon]
asn [geoip][asn] [<target>][as][number]
as_org [geoip][as_org] [<target>][as][organization][name]

target 的默认值也随模式变。disabled 下固定是 geoip;ECS 下没有固定默认值,而是从 source 推导:source 写成 [client][ip] 这种 ip 子字段时,target 自动取父字段 clientsourceclient_ip 这种平铺字段名时,target 变成必填项,不写会抛 ConfigurationError,Logstash 直接起不来。ECS 下 target 的合法取值是 client / destination / host / observer / server / source 六个,填别的不会拦你,但会在启动日志里留一条 warn。

最容易撞上的组合是 ECS 模式配 target => "geo":geo 子树会再嵌一层,实际路径变成 [geo][geo][country_iso_code]。按 disabled 时代的经验去 Kibana 里搜 geo.country_code2,什么都搜不到。示例配置里显式写清模式比依赖默认值省事,本文后面统一用 disabled;跑在 ECS 默认下的读者按上表把字段路径换一遍即可。

查不到 IP 时会发生什么

内网 IP 和保留地址段(RFC1918、loopback、link-local)在 MaxMind 库里没有记录。查不到时 geoip 不往 target 写任何东西,同时给 event 打上 _geoip_lookup_failuretag_on_failure 的默认值就是 ["_geoip_lookup_failure"])。用 City 库时还有一条更硬的规则:拿不到完整的经纬度对,整次 enrichment 直接中止——不是逐字段留空,而是一个字段都不写。

既然失败会留下 tag,就可以拿它继续分支,把内网流量和公网流量分开统计:

1
2
3
4
5
6
7
8
9
10
11
12
filter {
geoip {
source => "clientip"
ecs_compatibility => disabled
}

if "_geoip_lookup_failure" in [tags] {
mutate { add_field => { "traffic_scope" => "internal" } }
} else {
mutate { add_field => { "traffic_scope" => "external" } }
}
}

默认应该选这种写法,而不是在 geoip 之前用正则把内网段筛掉:正则要枚举三段 RFC1918 加 loopback 加 link-local,写错一个区间就静默漏判,而 tag 是查表的真实结果。只有当内网 IP 占了绝大多数、无效查表的开销已经上了 profile 时,才值得反过来用前置条件跳过 geoip。

数据库更新是有硬期限的运维任务

MaxMind 已把 GeoIP 库的授权从 Creative Commons 换成了 EULA,而 EULA 要求使用方在库更新后 30 天内跟上。Logstash 对此的实现方式,比"精度会随数据库新旧下降"严厉得多:

  • 自动更新由 GeoIP database download manager 负责,从 Elastic 的 endpoint 拉取,不是从 MaxMind 账户下载。开关是 logstash.yml 里的 xpack.geoip.downloader.enabled默认开启,每天检查一次。
  • 一旦 Logstash 切换到 EULA 库之后连续 30 天没能成功检查更新,geoip filter 会直接停止 enrich,并给 event 打上 _geoip_expired_database。这是为满足 EULA 合规的主动降级,不是 bug。
  • 状态可以从 Node Stats API 读:curl 'localhost:9600/_node/stats/geoip_download_manager',看 database.*.statusup_to_date / to_be_expired(25 天)/ expired(30 天))和 fail_check_in_days

对断网或气隙环境,这条规则的后果是:机器上线时地理字段一切正常,跑满一个月之后地理字段静默消失,而管道不报错、吞吐不变、日志里只多一个 tag。这跟"城市精度下降"完全不是一个量级的问题。

气隙环境的正解是把 xpack.geoip.downloader.enabled 设成 false。关掉之后 Logstash 会一直用捆绑的 CC 库,不受 30 天限制。代价是已经下载过的 EULA 库会被删除,以及从此手动管理更新。另一条路是自建 endpoint:用 Elasticsearch 自带的 bin/elasticsearch-geoip -s <目录> 把 MaxMind 下载来的 mmdb 生成静态服务文件,静态托管后配 xpack.geoip.download.endpoint 指过去,自动更新链路照旧工作。

至于精度本身:城市级定位在任何版本的库里都是估算,不要当作精确坐标使用。

条件块与 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" ecs_compatibility => disabled }
}
}
}

条件操作符覆盖:

1
2
3
4
5
==  !=  <  >  <=  >=          比较(字符串或数字)
=~ !~ 正则匹配 / 不匹配
in not in 成员检查(tag 数组、字段值)
and or xor nand 逻辑运算(二元,闭集就这四个)
! 取反(一元)

逻辑运算这一行最容易记错。官方列出的 boolean operator 是 and / or / xor / nand 四个,not 不在里面。not 只在两个位置出现——跟 in 组成 not in 这个成员检查运算符,以及作为一元取反 ! 的同义写法。想表达"两个条件都不成立",写法是 !([a] and [b]) 或者 nand,不能写 [a] not [b]。表达式可以用括号分组、可以嵌套。

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

tag 驱动分支让前序步骤的结果成为后续步骤的路由依据,配置里不必硬编码"哪类日志走哪条路"。这一点决定了 Logstash 的路由能力上限:所有分支判断都基于 event 自身携带的字段和 tag,没有外部路由表、没有订阅关系、没有 broker 侧的元数据。Kafka 用 topic 与 partition 在传输层就把消息分开,Logstash 则要等 event 进入 pipeline、经过前序 filter 加工出信号之后,才能在配置里做分流。好处是路由规则和加工逻辑写在一起、可以任意复杂;代价是分流发生在处理链内部,分不掉计算量。

内置 tag 由插件自动写入,自定义 tag 由 mutate 的 add_tag 添加,两者在条件里没有区别。

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

失败信号表

这套机制真正好用的前提是知道有哪些信号可以读。本文涉及的三个插件加上前两篇的 grok/dissect,失败信号是这一组:

信号 由谁写入 触发条件 event 去向
_grokparsefailure grok 所有 match pattern 全部匹配失败 字段没切出来,event 继续往下走
_groktimeout grok 单次匹配超过 timeout_millis(默认 30000) 匹配中止,event 继续往下走
_dissectfailure dissect 分隔符结构与 pattern 不符 字段没切出来,event 继续往下走
_dateparsefailure date match 里所有格式串都解析失败 @timestamp 保持进入时刻不变
_geoip_lookup_failure geoip IP 查不到,或 City 库拿不到经纬度对 target 下一个字段都不写
_geoip_expired_database geoip 连续 30 天未成功检查数据库更新 enrich 整体停止,地理字段全部消失
自定义 tag mutate add_tag 由条件块决定 由后续条件块决定

这张表里没有一行会中断 pipeline。filter 阶段的失败统一表现为"打个标记、继续往下走",要不要丢弃、要不要分流、要不要落到单独索引,都得自己写条件块。这跟 output 阶段完全不同。那边的失败会真的丢 event,见下一篇。

实验:一条 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
35
36
37
38
39
# 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"]
ecs_compatibility => disabled
}
}
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 已被删除;_grokparsefailure 不在 tags 里,if "_grokparsefailure" not in [tags] 这个条件成立,后续加工全部执行。

geo 字段那一格的结果会跟直觉相反。203.0.113.42 属于 RFC 5737 的 TEST-NET-3 文档保留段,MaxMind 库里没有记录,所以 geoip 查表失败,geo 字段根本不会出现tags 里多一个 _geoip_lookup_failure。换成真实公网 IP 重跑,才会看到 [geo][country_code2][geo][city_name]。这里的路径是两层扁平结构,是因为配置里显式写了 ecs_compatibility => disabled;跑在 8.x 默认的 ECS v8 模式下,同样的 target => "geo" 会把结果放到 [geo][geo][country_iso_code],字段名和层级都不一样。

模式提炼

1
2
3
4
5
6
7
8
9
10
11
12
13
14
模式:tag 驱动的渐进式加工链

- 前序 filter 在 tag 里写入结果信号(成功/失败/类型),
后续 filter 读 tag 决定是否执行,形成懒求值的加工链
- date filter 改写 @timestamp 决定 event 落在哪条时间轴上,
时序查询与 ILM 都依赖它,属于必做项
- 字段类型转换有两条路径:grok 内联后缀覆盖 int/float,
mutate.convert 额外覆盖 boolean 与欧洲数字格式
- geoip 在进程内查 mmdb,延迟来自内存查找,成本随库大小上升
- 条件块是 pipeline 内路由,不等于 pipeline 间路由
- filter 阶段的失败统一表现为"打 tag 后继续往下走",
丢弃、分流、告警都得自己写条件块
- 插件的字段落位受 ecs_compatibility 支配,
示例配置显式写出模式比依赖默认值可靠

这套"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 做的是同一件事。两者都在"数据入流"之后、"时序计算"之前,把处理时间替换成事件时间——这是所有流处理系统里时序正确性的共同前提。

常见误解

误解一:“filter 里的多个 mutate 操作可以合并成一个”。mutate 内部的操作有固定执行顺序,与书写顺序无关。源码 mutate.rbfilter 方法就是一串顺序调用,共 14 步:

1
2
3
coerce → rename → update → replace → convert → gsub
→ uppercase → capitalize → lowercase → strip
→ split → join → merge → copy

add_tag / remove_tag / add_field / remove_field 这四个不在这张表里。它们是所有 filter 插件共有的 common options,由基类的 filter_matched(event) 统一处理,而这个调用在 mutate 主体 14 步全部跑完之后才发生。所以这四个操作永远最后生效,把它们排进同一张有序表会误导对机制的理解。filter_matched 内部的顺序是 add_fieldremove_fieldadd_tagremove_tag

官方文档对此挂了一条 IMPORTANT 级约束:

Each mutation must be in its own code block if the sequence of operations needs to be preserved.

顺序敏感的场景就拆块,别指望在一个块里调整书写顺序。举两个例子:先 rename 再基于新名字 gsub,写在同一块里能按预期工作,因为 rename 在第 2 步、gsub 在第 6 步;但想先 add_fieldrename 这个新字段,同一块里做不到,因为 rename 在主体里跑、add_fieldfilter_matched 里跑,实际顺序永远是 rename 在前。本文开头 mutate 示例里 replaceremove_field 抹掉,也是同一条规则的后果。

误解二:“geoip 数据库配好就能一直用下去”。装完确实能用,但有个 30 天的隐形期限。Logstash 默认开启 GeoIP database download manager,每天检查一次更新;一旦切换到 EULA 库之后连续 30 天检查失败,geoip 会为满足 MaxMind EULA 合规而停止 enrich,只给 event 打 _geoip_expired_database。断网环境里这表现为地理字段在上线一个月后静默消失,而管道不报错。气隙部署要显式把 xpack.geoip.downloader.enabled 设成 false,让它一直用捆绑的 CC 库。

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

误解四:“geoip 会拖慢管道,因为要查外部服务”。geoip 数据库在 Logstash 启动时加载进内存,查询完全在进程内完成,没有网络 IO。性能瓶颈是内存读取和 MaxMind 数据库的查找开销,与外部服务无关。数据库文件越大,冷启动加载时间越长,但单次查询延迟通常在微秒级。插件还带一层 cache_size(默认 1000)的查表缓存,日志里相邻行 IP 重复率高时命中率不低。注意它没有淘汰策略,缓存满了就不再新增条目。

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

练习

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

  2. 在一个 mutate 块里同时写 add_field => { "tmp" => "hello" }rename => { "tmp" => "final" },观察 rubydebug 输出里剩下的是 tmp 还是 final。再把两个操作拆成前后两个 mutate 块重跑一次。结合 14 步顺序表和 filter_matched 的位置解释两次结果为什么不同,并说明为什么换成 rename + gsub 这一对时,写在一个块里反而没问题。

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

  4. 把实验那份配置里 geoip 的 ecs_compatibility => disabled 删掉(让它回落到 8.x 默认的 v8),换一个真实公网 IP 重跑,对比两次 rubydebug 输出里地理字段的完整路径。再把 sourceclientip 改成 [client][ip] 并删掉 target,观察 target 是否按父字段自动推导;然后把 source 改回 clientip 且仍不写 target,确认 Logstash 会不会启动失败、报什么错。

系列导航

序号 主题
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 的冲击

参考资料