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

