You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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)**功能,实时生成聚合后的索引。

操作步骤

  1. 先让Logstash把所有原始数据输出到一个原始索引(比如job-logs-raw)
  2. 在ES/OS中创建一个Transform:
    • 数据源选择原始索引
    • 分组字段选job_id
    • 聚合规则设置为“将所有事件的字段收集到数组”(比如用scripted_metric或top_hits实现,具体根据需求调整)
    • 设置同步模式为实时同步,新的原始数据进来时,聚合索引会自动更新

优点

  • 原始数据完全保留,聚合逻辑可以随时调整,不需要修改Logstash配置
  • 不用担心Logstash的内存问题,聚合计算由ES/OS集群承担

方案三:用Logstash的multiline+后续处理(适合简单场景)

如果同一个job的事件是连续输出的,也可以先用multiline过滤器把同job的事件合并成单条,再同时输出原始和合并后的数据。但这个方法只适合事件顺序严格的场景,容错性不如前两个方案。


备注:内容来源于stack exchange,提问作者fmakawa

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.16 10:22:58