如何在Logstash中同时保留原始日志数据与聚合后的数据(或借助Elastic/OpenSearch实现聚合)
如何在Logstash中同时保留原始日志数据与聚合后的数据(或借助Elastic/OpenSearch实现聚合)
嘿,这个需求其实挺常见的,我给你几个靠谱的方案,看看哪个更贴合你的场景:
方案一:在Logstash中分流处理——同时保存原始数据+生成聚合数据
这个方法直接在Logstash里做事件复制,一份走原始输出,另一份做聚合后输出到不同索引,完全满足你“不丢原始数据”的要求。
核心思路是用clone过滤器复制事件,给不同分支打标签区分,然后分别处理:
完整配置示例
input { tcp { port => 4512 codec => json } } filter { # 给所有原始事件打标记,用来识别原始流 mutate { add_tag => ["raw_event"] } # 复制一份事件,专门用来做聚合 clone { clones => ["aggregation_event"] } # 只对聚合分支的事件做聚合处理 if "aggregation_event" in [tags] { aggregate { # 按job_id分组 task_id => "%{job_id}" # 把每个事件的内容收集到数组里(可根据实际字段调整) code => " map['job_id'] = event.get('job_id') map['job_events'] ||= [] map['job_events'] << event.to_hash # 如果是job结束的事件(假设你有标识字段,比如event_type为job_finished),标记为完成 if event.get('event_type') == 'job_finished' event.set('aggregation_complete', true) end " # 当标记了完成时,触发聚合结果输出 flush_on_timeout => false flush_when => "aggregation_complete" # 如果job没有结束标识,可设置超时自动flush(根据业务调整时长) timeout => 3600 } # 聚合完成后,移除多余标签避免干扰输出 mutate { remove_tag => ["aggregation_event"] } } } output { # 原始事件输出到原始索引 if "raw_event" in [tags] { elasticsearch { hosts => ["your-es-host:9200"] index => "job-logs-raw-%{+YYYY.MM.dd}" } } # 聚合后的事件输出到聚合索引 if "aggregation_complete" in [fields] { elasticsearch { hosts => ["your-es-host:9200"] index => "job-logs-aggregated-%{+YYYY.MM.dd}" } } }
注意点
- 要根据实际业务调整聚合逻辑:比如收集哪些字段、用什么标识job结束(如果没有结束标识,就靠
timeout自动触发) - 如果job运行时间长、数据量大,建议开启Logstash的
persistent_queue,避免重启丢失聚合中的数据
方案二:在Elastic/OpenSearch端做后聚合——只存原始数据,自动生成聚合索引
如果不想在Logstash里做复杂的聚合逻辑,也可以把所有原始数据先存入ES/OS,然后用官方的**Transform(Elastic)或Index Transforms(OpenSearch)**功能,实时生成聚合后的索引。
操作步骤
- 先让Logstash把所有原始数据输出到一个原始索引(比如
job-logs-raw) - 在ES/OS中创建一个Transform:
- 数据源选择原始索引
- 分组字段选
job_id - 聚合规则设置为“将所有事件的字段收集到数组”(比如用
scripted_metric或top_hits实现,具体根据需求调整) - 设置同步模式为实时同步,新的原始数据进来时,聚合索引会自动更新
优点
- 原始数据完全保留,聚合逻辑可以随时调整,不需要修改Logstash配置
- 不用担心Logstash的内存问题,聚合计算由ES/OS集群承担
方案三:用Logstash的multiline+后续处理(适合简单场景)
如果同一个job的事件是连续输出的,也可以先用multiline过滤器把同job的事件合并成单条,再同时输出原始和合并后的数据。但这个方法只适合事件顺序严格的场景,容错性不如前两个方案。
备注:内容来源于stack exchange,提问作者fmakawa
相关产品推荐
相关产品推荐

