Logstash对接Elasticsearch:清除旧数据与批量事件处理配置
嘿,针对你问到的两个Logstash+Elasticsearch的问题,再结合你补充的数组拆分需求,我给你整理了实际可落地的解决方案:
问题1:Logstash输出新数据前清除所有旧数据
这里分两种常用场景,你可以根据自己的情况选择:
场景1:每次运行Logstash前清空整个目标索引
最简单的方式是在启动Logstash前,用脚本调用Elasticsearch的删除API来清空索引。比如写个shell脚本:# 先删除目标索引 curl -X DELETE "http://你的ES地址:9200/目标索引名" # 再启动Logstash /path/to/logstash -f your-config.conf这种方式直观且不容易出错,适合全量替换数据的场景。
场景2:在Logstash流程内触发全量删除
如果想在Logstash处理过程中自动执行删除,可以借助elasticsearch输出插件的delete_by_query动作。不过要注意只触发一次,避免重复删除:
在filter阶段生成一个删除触发事件:filter { # 仅在启动时生成一次删除事件(可以通过input的schedule或其他条件控制) if [init_trigger] == "delete_all" { mutate { add_field => { "[@metadata][_op_type]" => "delete_by_query" } add_field => { "[@metadata][_index]" => "目标索引名" } add_field => { "[@metadata][query]" => '{"match_all": {}}' } } } }然后输出配置里处理这个删除事件:
output { if [init_trigger] == "delete_all" { elasticsearch { hosts => ["你的ES地址:9200"] action => "%{[@metadata][_op_type]}" index => "%{[@metadata][_index]}" query => "%{[@metadata][query]}" } } else { # 正常输出新数据的配置 elasticsearch { hosts => ["你的ES地址:9200"] index => "目标索引名" } } }
问题2:每周获取多组事件并删除对应旧事件
结合你补充的拆分数组为单个文档的需求,我把方案拆成两步:先处理数组拆分,再实现旧数据删除。
第一步:拆分records数组为单个文档
因为你提到内部对象没有固定schema,不能用Nested Object,所以用Logstash的split插件刚好合适,配置如下:
filter { # 把records数组的每个元素拆成独立文档,放到record字段里 split { field => "records" target => "record" remove_field => ["records"] # 拆分后移除原数组字段 } }
这样输入的{host:"host1", type:"packages", records: [...]}就会被拆成多个符合你期望的{host:"host1", type:"packages", record: {...}}文档。
第二步:每周更新时删除对应host+type的旧事件
因为你是每周拉取全量数据,最合适的方式是先删除该主机对应类型的所有旧文档,再插入新的拆分文档。配置如下:
Filter阶段:生成删除事件
filter { # 复制原始事件(拆分前)生成删除请求事件 clone { clones => ["delete_old_event"] add_tag => ["delete_old"] } # 配置删除事件的动作和查询条件(匹配当前host和type) if "delete_old" in [tags] { mutate { add_field => { "[@metadata][_op_type]" => "delete_by_query" } add_field => { "[@metadata][_index]" => "目标索引名" } # 精准匹配当前host和type的所有旧文档 add_field => { "[@metadata][query]" => '{"bool": {"must": [{"term": {"host": "%{host}"}}, {"term": {"type": "%{type}"}}]}}' } remove_field => ["records"] # 删除事件不需要这个字段 } } # 拆分数组(只处理非删除事件) if "delete_old" not in [tags] { split { field => "records" target => "record" remove_field => ["records"] } } }
Output阶段:分别处理删除和插入
output { # 执行删除旧文档的请求 if "delete_old" in [tags] { elasticsearch { hosts => ["你的ES地址:9200"] action => "%{[@metadata][_op_type]}" index => "%{[@metadata][_index]}" query => "%{[@metadata][query]}" } } # 插入拆分后的新文档 else { elasticsearch { hosts => ["你的ES地址:9200"] index => "目标索引名" action => "index" } } }
注意事项
- 要确保Elasticsearch的账号拥有
delete_by_query的权限 - 如果数据量较大,建议在业务低峰期执行每周的同步任务,避免影响ES性能
- 可以通过Logstash input的
schedule配置(比如schedule => "0 0 * * 0"表示每周日凌晨)来触发每周的拉取和更新
内容的提问来源于stack exchange,提问作者osexp2000
相关产品推荐
相关产品推荐

