使用Elasticsearch Pipeline去除NDJSON数据中的空数组
方案一:使用Elasticsearch Ingest Pipeline处理
由于tagset.username下的键是动态变化的(每条日志对应不同域账号),grok无法处理这类动态结构,需用script处理器提取动态键:
- 创建Pipeline配置
PUT _ingest/pipeline/extract-username { "description": "从tagset.username的动态键中提取域账号", "processors": [ { "script": { "source": """ def usernameMap = ctx.tagset?.username; if (usernameMap != null && usernameMap instanceof Map) { for (def entry : usernameMap.entrySet()) { ctx.tagset.username = entry.getKey(); break; } } """ } }, // 可选:将username移至根字段 { "rename": { "field": "tagset.username", "target_field": "username", "ignore_missing": true } } ] }
- 配置Filebeat指定该Pipeline
修改filebeat.yml的Elasticsearch输出段:
output.elasticsearch: hosts: ["你的ES地址:9200"] pipeline: "extract-username"
方案二:使用Logstash处理
Logstash的ruby过滤器可灵活处理动态键结构:
- Logstash配置文件示例
input { beats { port => 5044 } } filter { ruby { code => """ username_map = event.get('[tagset][username]') if username_map.is_a?(Hash) && !username_map.empty? domain_username = username_map.keys.first event.set('[tagset][username]', domain_username) # 可选:将账号写入根字段username event.set('username', domain_username) end """ } } output { elasticsearch { hosts => ["你的ES地址:9200"] index => "你的索引名" } }
- 修改Filebeat输出至Logstash
修改filebeat.yml:
output.logstash: hosts: ["你的Logstash地址:5044"]
测试提示
- Pipeline可通过
_ingest/pipeline/extract-username/_simulate接口传入测试日志验证效果 - 确保所有组件版本(Filebeat/Elasticsearch/Logstash)均为7.3,避免兼容性问题
内容的提问来源于stack exchange,提问作者Mikolaj
相关产品推荐
相关产品推荐

