You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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
  }
}

配置说明

  1. 字段提取:通过add_field将_source嵌套结构内的目标字段复制到根层级,使用Logstash的字段引用语法%{[字段路径]}实现取值。
  2. 字段清理:通过remove_field直接删除所有不需要的字段,包括根层级的_id、_score,以及整个_source字段(连带其内部多余字段),确保最终输出仅保留所需内容。

内容的提问来源于stack exchange,提问作者Martin Dvoracek

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 01:45:31