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
相关产品推荐
相关产品推荐

