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

Logstash对接Kafka与Elasticsearch时if/else逻辑异常求助

你的Logstash分支逻辑问题排查与修复

首先,你的输出配置存在语法错误,这是导致数据无法写入Elasticsearch的核心原因,我帮你拆解问题并给出修正方案:

1. 核心语法错误

在else分支里,你错误地直接嵌套了条件判断:

else { [kafka][topic_id] == "newdata" { ... } }

Logstash的配置语法中,else后面不能直接跟条件表达式,多个条件分支需要用else if来串联,而不是在else块里再写判断逻辑。

2. 字段判断的错误

当你开启decorate_events => true时,Kafka的主题信息会被放到[kafka][topic]字段中(而非topic_id);另外,你的输入配置已经订阅了明确的两个主题,直接判断顶层的[topic]字段也完全可行。

3. 修正后的完整配置

下面是修复后的输出配置,同时保持了你的业务逻辑——为两个主题创建独立索引:

input {
  kafka {
    bootstrap_servers => '192.168.2.66:9092,192.168.2.67:9092,192.168.2.68:9092'
    topics => ["linum", "newdata"]
    codec => "json"
    type => "kafka"
    tags => ["linumsrvr1"]
    decorate_events => true
  }
}

output {
  # 处理linum主题的数据
  if [topic] == "linum" {
    elasticsearch {
      hosts => ["192.168.2.32:9200"]
      action => "index"
      index => "linum-%{+YYYY.MM.dd}"
    }
    stdout { codec => rubydebug }
  }
  # 处理newdata主题的数据
  else if [topic] == "newdata" {
    elasticsearch {
      hosts => ["192.168.2.32:9200"]
      index => "newdata-%{+YYYY.MM.dd}"
    }
    stdout { codec => rubydebug }
  }
}

额外排查建议

如果替换配置后仍有问题,建议查看Logstash的日志文件(默认路径为/var/log/logstash/logstash-plain.log),Logstash会在日志中明确标出配置语法错误、字段不存在等问题,这是快速定位问题的重要依据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:16:12