如何在Logstash中过滤Kafka消息,仅保留指定字段输出
Logstash过滤Kafka输入消息,保留指定字段
问题场景
从Kafka输入到Logstash的消息结构如下:
{ "_index" : "progress", "_id" : "Q27Y2IYBIZUq2eJ6WJQR", "_score" : 1.0, "_source" : { "itemId" : 3, "weight" : 358, "timeStartedInMillis" : 39131, "flow" : 3725, "@version" : "1", "@timestamp" : "2023-03-13T02:41:42.313784Z", "event" : { "original" : "{\"timeStartedInMillis\": 39131, \"procedureId\": 3, \"temperature\": 47, \"weight\": 358, \"flow\": 3725}" }, "type" : "log", "temperature" : 47, "tags" : [ "kafka_source" ] } }
需要过滤后仅输出以下格式:
{ "_index" : "progress", "itemId" : 3, "weight" : 358, "timeStartedInMillis" : 39131, "flow" : 3725, "temperature" : 47 }
解决方案
使用Logstash的mutate插件完成字段提取和清理,配置示例如下:
Filter阶段配置
filter { # 将_source中的目标字段提取到根层级 mutate { add_field => { "itemId" => "%{[ _source ][ itemId ]}", "weight" => "%{[ _source ][ weight ]}", "timeStartedInMillis" => "%{[ _source ][ timeStartedInMillis ]}", "flow" => "%{[ _source ][ flow ]}", "temperature" => "%{[ _source ][ temperature ]}" } # 删除所有不需要的字段 remove_field => [ "_id", "_score", "_source", "@version", "@timestamp", "event", "type", "tags" ] } }
Output阶段配置(以标准输出为例)
output { stdout { codec => json_lines } }
配置说明
- 字段提取:通过
add_field将_source嵌套结构内的目标字段复制到根层级,使用Logstash的字段引用语法%{[字段路径]}实现取值。 - 字段清理:通过
remove_field直接删除所有不需要的字段,包括根层级的_id、_score,以及整个_source字段(连带其内部多余字段),确保最终输出仅保留所需内容。
内容的提问来源于stack exchange,提问作者Martin Dvoracek
相关产品推荐
相关产品推荐

