深入 Logstash 07 - 常用 filter 组合:mutate、date、geoip 与条件
上一篇讲完了 dissect 如何用分隔符驱动的线性扫描替代 Grok 的回溯正则,代价是无法处理非均匀格式。这一篇进入 filter 阶段最常见的后续加工链:mutate 改写字段、date 把日志时间戳覆盖 @timestamp、geoip 用 IP 换取地理信息,以及用条件判断和 tag 做分支路由。
核心问题:date filter 为什么一定要改写 @timestamp,而不是新增一个字段;条件判断和 tag 驱动的分支是 Logstash 唯一的路由原语,它和消息系统的 topic 路由有何本质区别。
数据流:event 经过 filter 链的变换过程
1 | |
filter block 是一个顺序执行的指令序列。所有 filter 插件和条件块从上到下依次运行,每一步的输出是下一步的输入。没有隐式的并行——并行发生在 pipeline worker 层面,而不在单个 event 的处理链里。
关键对象
filter 阶段操作的对象始终是 event 的字段映射(field map)。每个 filter 插件本质上是一个 event.fields → event.fields 的函数,读取已有字段,写入或删除字段。
四个常用插件各自的职责:
1 | |
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 | |
这段配置故意保留了一个顺序冲突:replace 往 message 里写了 processed: ...,remove_field 又把 message 删掉,两个操作还在同一个 mutate 块里。结果是输出里连 message 字段都不存在,那次 replace 完全白做。想保住 replace 的结果,就得把 remove_field 拆到后面一个独立的 mutate 块里。为什么是这个顺序,下文"误解一"里展开。
Grok 切出来的字段默认全是字符串,写进 ES 之后如果映射成了 keyword,数值聚合会失败。把字符串变成数字有两条路径:一条是 grok 的内联类型后缀 %{NUMBER:response:int},在切字段的同时就完成转换;另一条是 mutate 的 convert,在切完之后单独转一遍。内联后缀省一步,但官方写明它只支持 int 和 float 两种转换。需要 boolean,或者要解析欧洲数字格式(1.234,56 这种拿逗号当小数点的写法,对应 integer_eu / float_eu),就只有 convert 能做。这是它无法被内联后缀替代的地方。
gsub 接受三元组列表:[字段名, 正则, 替换字符串],可以连续列多组。它在字段内容里做全局正则替换,常见用途是从 URL 里去掉查询参数、清洗特殊字符。
remove_field 在写 ES 之前清理掉 message(原始文本)和中间字段,能显著降低单文档大小和 ES 的 _source 存储压力。
date:@timestamp 为何必须被改写
1 | |
默认情况下,@timestamp 是 event 进入 Logstash 的时刻,不是日志里那条记录真正发生的时刻。两者之间的差距取决于管道的延迟——正常情况下差几秒,批量补录历史日志时可能差几小时甚至几天。
如果不改写 @timestamp,Kibana 的时间线查询(range query on @timestamp)会按处理时间而非事件时间过滤,导致历史日志全部堆在处理时刻附近、事件的真实时序完全丢失。Elasticsearch 的 rollover/ILM 按 @timestamp 决定文档落在哪个索引分片,用处理时间同样会破坏基于时间的索引分布。
match 数组的第一个元素是存放时间字符串的字段名,后面跟一个或多个格式串。多个格式串是或关系——Logstash 依次尝试,第一个匹配成功就用,全部失败则打 _dateparsefailure tag。格式串使用 Joda-Time 语法(Logstash 运行在 JVM 上,时间解析底层是 Java 库)。
除了自定义格式串,match 还接受四个内置的格式字面量,覆盖了实际日志里最常见的几种时间形态:
1 | |
这四个里前两个在实际日志里出现频率最高。结构化日志框架(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 | |
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 自动取父字段 client;source 是 client_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_failure(tag_on_failure 的默认值就是 ["_geoip_lookup_failure"])。用 City 库时还有一条更硬的规则:拿不到完整的经纬度对,整次 enrichment 直接中止——不是逐字段留空,而是一个字段都不写。
既然失败会留下 tag,就可以拿它继续分支,把内网流量和公网流量分开统计:
1 | |
默认应该选这种写法,而不是在 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.*.status(up_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 | |
条件操作符覆盖:
1 | |
逻辑运算这一行最容易记错。官方列出的 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 | |
输入一行:
1 | |
rubydebug 输出中可以观察到:@timestamp 已从进入时刻改写为 2025-08-10T13:55:36.000Z;response 和 bytes 是整数;log_ts 和 message 已被删除;_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 | |
这套"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.rb 的 filter 方法就是一串顺序调用,共 14 步:
1 | |
add_tag / remove_tag / add_field / remove_field 这四个不在这张表里。它们是所有 filter 插件共有的 common options,由基类的 filter_matched(event) 统一处理,而这个调用在 mutate 主体 14 步全部跑完之后才发生。所以这四个操作永远最后生效,把它们排进同一张有序表会误导对机制的理解。filter_matched 内部的顺序是 add_field → remove_field → add_tag → remove_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_field 再 rename 这个新字段,同一块里做不到,因为 rename 在主体里跑、add_field 在 filter_matched 里跑,实际顺序永远是 rename 在前。本文开头 mutate 示例里 replace 被 remove_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 实验确认字段是否存在。
练习
-
用 rubydebug 管道验证 date filter 改写效果:构造一条带历史时间戳的日志(比如一个月前的时间),分别观察有 date filter 和没有 date filter 时
@timestamp的差异。确认 Kibana 的时间线查询在两种情况下会落在哪个时间段。 -
在一个 mutate 块里同时写
add_field => { "tmp" => "hello" }和rename => { "tmp" => "final" },观察 rubydebug 输出里剩下的是tmp还是final。再把两个操作拆成前后两个 mutate 块重跑一次。结合 14 步顺序表和filter_matched的位置解释两次结果为什么不同,并说明为什么换成rename+gsub这一对时,写在一个块里反而没问题。 -
构造一条解析失败的日志(让 Grok 匹配不上),观察
tags数组里出现_grokparsefailure;然后在 filter 里加条件块,对_grokparsefailure打add_tag => ["needs_manual_review"],并add_field => { "parse_status" => "failed" }。验证 rubydebug 输出里两个 tag 都存在。 -
把实验那份配置里 geoip 的
ecs_compatibility => disabled删掉(让它回落到 8.x 默认的v8),换一个真实公网 IP 重跑,对比两次 rubydebug 输出里地理字段的完整路径。再把source从clientip改成[client][ip]并删掉target,观察target是否按父字段自动推导;然后把source改回clientip且仍不写target,确认 Logstash 会不会启动失败、报什么错。
系列导航
参考资料
- Logstash mutate filter 文档:https://www.elastic.co/guide/en/logstash/current/plugins-filters-mutate.html(操作列表与执行顺序)
- Logstash date filter 文档:https://www.elastic.co/guide/en/logstash/current/plugins-filters-date.html(match 格式串与 Joda-Time 语法)
- Logstash geoip filter 文档:https://www.elastic.co/guide/en/logstash/current/plugins-filters-geoip.html(数据库配置与字段列表)
- Logstash 条件判断文档:https://www.elastic.co/guide/en/logstash/current/event-dependent-configuration.html(比较、boolean、一元操作符的完整闭集与字段引用语法)
- Logstash ECS 兼容模式:https://www.elastic.co/guide/en/logstash/current/ecs-ls.html(
pipeline.ecs_compatibility默认值与插件级、pipeline 级覆盖方式) - mutate filter 源码:https://github.com/logstash-plugins/logstash-filter-mutate/blob/main/lib/logstash/filters/mutate.rb(
filter方法就是那 14 步顺序调用) - filter 基类源码:https://github.com/elastic/logstash/blob/main/logstash-core/lib/logstash/filters/base.rb(
filter_matched里 add_field/remove_field/add_tag/remove_tag 的统一处理) - Joda-Time format 参考:https://www.joda.org/joda-time/apidocs/org/joda/time/format/DateTimeFormat.html(date filter 格式串底层实现)
- MaxMind GeoLite2 数据库:https://dev.maxmind.com/geoip/geolite2-free-geolocation-data(geoip 数据库来源与更新策略)
