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

Logstash按消息ENV值分发至不同Kafka Topic时消息丢失问题排查

问题排查:Logstash按ENV字段路由Kafka Topic时消息丢失

原始正常配置

output {
  kafka {
    topic_id => "topic1"
    bootstrap_servers => "172.172.172.172:1722"
    codec => json
    acks => "0"
    partitioner => round_robin
    compression_type => gzip
    max_request_size => 10485760
  }
}

修改后的条件路由配置(出现消息丢失)

output {
  if [ENV] == "DEV" {
    kafka {
      topic_id => "topic_DEV"
      bootstrap_servers => "172.172.172.172:1722"
      codec => json
      acks => "0"
      partitioner => round_robin
      compression_type => gzip
      max_request_size => 10485760
    }
  }
  else if [ENV] == "PRD" {
    kafka {
      topic_id => "topic_Production"
      bootstrap_servers => "172.172.172.172:1722"
      codec => json
      acks => "0"
      partitioner => round_robin
      compression_type => gzip
      max_request_size => 10485760
    }
  }
}

排查方向及解决方案

1. 确认字段是否被正确解析到顶层

如果输入消息是JSON格式,但Logstash的input阶段未配置codec => json解析,ENV字段会被包裹在message字段的字符串中,而非顶层字段。此时[ENV]不存在,所有消息都无法匹配条件,直接被丢弃。

解决:在input模块添加JSON解析,示例:

input {
  file {
    path => "/path/to/your/logs"
    codec => json # 关键:将JSON消息解析为顶层字段
  }
}

2. 检查字段名大小写匹配

Logstash字段引用大小写敏感,确认消息中的字段是ENV(全大写)还是小写env。若实际字段为env,[ENV]无法匹配,导致条件不触发。

3. 补充兜底分支排查异常

即使确认只有DEV和PRD两个值,也建议添加else分支捕获未匹配消息,排查是否存在字段为空、值含空格等异常:

output {
  if [ENV] == "DEV" {
    kafka {
      topic_id => "topic_DEV"
      bootstrap_servers => "172.172.172.172:1722"
      codec => json
      acks => "0"
      partitioner => round_robin
      compression_type => gzip
      max_request_size => 10485760
    }
  }
  else if [ENV] == "PRD" {
    kafka {
      topic_id => "topic_Production"
      bootstrap_servers => "172.172.172.172:1722"
      codec => json
      acks => "0"
      partitioner => round_robin
      compression_type => gzip
      max_request_size => 10485760
    }
  }
  else {
    # 输出到控制台查看未匹配消息
    stdout { codec => rubydebug }
  }
}

4. 通过Logstash日志定位问题

开启Logstash DEBUG级日志,查看消息字段结构和条件匹配情况:

  • 修改log4j2.properties将日志级别设为DEBUG
  • 启动Logstash后,观察日志中condition match相关记录,确认ENV字段是否存在

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 12:05:58