Logstash输出失败时如何将事件转发至其他存储位置?
解决Logstash输出Elasticsearch失败事件的丢失问题
刚好碰到过类似的场景,给你两个实用的方案来处理这些因严格映射导致的400错误事件,确保数据不会丢失:
方案1:在输出阶段直接路由失败事件
利用Logstash elasticsearch输出插件的track_status参数,追踪输出失败的HTTP状态码,然后通过条件判断把失败事件转发到专门的存储位置(比如备用ES索引或文本日志)。
修改你的管道输出部分如下:
output { # 尝试将事件发送到主索引 elasticsearch { hosts => ["elasticsearch:9200"] index => "events-dev6-test" document_type => "_doc" manage_template => false # 开启状态追踪,失败事件会自动添加_elasticsearch_status字段 track_status => true } # 处理400错误的失败事件 if [_elasticsearch_status] == 400 { # 发送到专门的失败事件索引(建议给这个索引设置宽松映射,避免二次失败) elasticsearch { hosts => ["elasticsearch:9200"] index => "events-dev6-failed" document_type => "_doc" manage_template => false } # 同时输出到本地日志文件,方便后续排查问题 file { path => "/var/log/logstash/failed_events.log" codec => json_lines } } stdout { codec => rubydebug } }
方案说明
track_status => true会让Logstash在事件输出失败时,为事件添加_elasticsearch_status字段,记录对应的HTTP错误码(比如你的场景里的400)。- 我们通过
if [_elasticsearch_status] == 400的条件判断,把这些不符合映射的事件分流到备用索引和本地日志,既保证了数据不丢失,也能单独排查这些异常事件的结构问题。
方案2:启用死信队列(DLQ)持久化失败事件
如果你的场景中存在多次重试仍失败的事件,或者希望不影响主管道的性能,可以启用Logstash的死信队列,把持久化的失败事件单独处理。
步骤1:配置Logstash启用死信队列
修改Logstash的主配置文件logstash.yml,添加以下配置:
dead_letter_queue.enable: true dead_letter_queue.max_bytes: 1gb # 可根据你的存储需求调整大小
步骤2:创建单独的管道消费死信队列
新建一个管道配置文件(比如dlq-handler.conf),用来读取死信队列里的事件并转发:
input { dead_letter_queue { path => "/var/lib/logstash/dead_letter_queue" # 对应Logstash的DLQ存储路径 pipeline_id => "main" # 替换成你主管道的ID,默认是main start_position => "beginning" } } filter { # 给死信队列的事件添加标记,方便后续识别 mutate { add_field => { "[metadata][source]" => "dead_letter_queue" } } } output { # 发送到死信专用索引 elasticsearch { hosts => ["elasticsearch:9200"] index => "events-dev6-dlq-failed" document_type => "_doc" manage_template => false } # 或者输出到本地文件 file { path => "/var/log/logstash/dlq_failed_events.log" codec => json_lines } }
方案说明
- 死信队列会把主管道中多次重试仍失败的事件持久化到磁盘,不会丢失。
- 通过单独的管道处理DLQ事件,不会影响主管道的正常运行,且可以灵活调整处理逻辑。
额外小技巧:提前过滤异常事件
你还可以在filter阶段提前识别不符合结构的事件,提前标记并分流,减少ES的报错次数:
filter { mutate { remove_field => ["sequence", "host", "@version"] } # 提前检查事件结构,标记潜在的失败事件 if [name] == "wrong data" or (![data][query]) { mutate { add_tag => ["potential_failure"] } } }
然后在output阶段直接分流带potential_failure标签的事件:
if "potential_failure" in [tags] { elasticsearch { hosts => ["elasticsearch:9200"] index => "events-dev6-failed" document_type => "_doc" manage_template => false } }
内容的提问来源于stack exchange,提问作者Thomas Decaux
相关产品推荐
相关产品推荐

