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

如何在Elasticsearch中基于Jaeger的插入文档生成聚合新文档?

解决方案:实现Jaeger Span的聚合报告文档生成

你提到想用CDC但Elasticsearch不支持原生CDC,这里有几个可行的方案来完成你的需求,从简单到复杂依次介绍:

方案1:使用Elasticsearch Watcher(原生告警功能)

这是最轻量化的方案,利用Elasticsearch自带的Watcher功能来监听新插入的phase=end文档,并自动完成聚合逻辑。

步骤说明

  1. 创建Watcher:配置一个Watcher,监听jaeger-span-*索引的文档新增事件,过滤出带有tags.key:phase且tags.value:end的文档。
  2. 查询关联的初始Span:在Watcher的执行阶段,通过Search查询同一个traceID下的所有Span,找到最小的startTime值(也就是初始文档的startTime)。
  3. 计算聚合字段:
    • 计算endTime:当前phase=end文档的startTime + duration
    • 计算总duration:endTime - 初始Span的startTime
  4. 写入聚合文档:使用Watcher的index动作,将计算后的聚合数据写入指定的报告索引(比如jaeger-aggregated-reports)。

简化的Watcher DSL示例

{
  "trigger": {
    "schedule": {
      "interval": "1m" // 每分钟扫描一次新增文档,可根据需求调整
    }
  },
  "input": {
    "search": {
      "request": {
        "indices": ["jaeger-span-*"],
        "body": {
          "query": {
            "bool": {
              "must": [
                {"term": {"tags.key": "phase"}},
                {"term": {"tags.value": "end"}}
              ],
              "filter": {
                "range": {
                  "@timestamp": { // 只扫描最近1分钟的新增文档,避免重复处理
                    "gte": "now-1m"
                  }
                }
              }
            }
          }
        }
      }
    }
  },
  "condition": {
    "compare": {
      "ctx.payload.hits.total": {
        "gt": 0
      }
    }
  },
  "actions": {
    "aggregate_and_index": {
      "foreach": "ctx.payload.hits.hits",
      "max_iterations": 100,
      "action": {
        "chain": {
          "actions": [
            // 步骤1:查询当前traceID的所有Span,获取最小startTime
            {
              "search": {
                "request": {
                  "indices": ["jaeger-span-*"],
                  "body": {
                    "query": {
                      "term": {"traceID": "{{_source.fields.traceID.0}}"}
                    },
                    "aggs": {
                      "min_start_time": {
                        "min": {
                          "field": "fields.startTime"
                        }
                      }
                    },
                    "size": 0
                  }
                }
              }
            },
            // 步骤2:计算聚合字段并写入报告索引
            {
              "index": {
                "index": "jaeger-aggregated-reports",
                "document_id": "{{_source.fields.traceID.0}}",
                "body": {
                  "fields": {
                    "traceID": "{{_source.fields.traceID}}",
                    "startTime": "{{ctx.payload.aggregations.min_start_time.value}}",
                    "endTime": "{{_source.fields.startTime.0 + _source.fields.duration.0}}",
                    "duration": "{{(_source.fields.startTime.0 + _source.fields.duration.0) - ctx.payload.aggregations.min_start_time.value}}",
                    "process.serviceName": "{{_source.fields.process.serviceName}}"
                  }
                }
              }
            }
          ]
        }
      }
    }
  }
}

注意事项

  • 要确保Watcher有足够的权限查询Jaeger索引并写入报告索引。
  • 可以添加一个过滤条件,避免重复处理同一个traceID(比如检查报告索引中是否已存在该traceID的文档)。
  • Jaeger的索引按天命名,所以查询时要使用通配符jaeger-span-*覆盖所有相关索引。

方案2:使用Logstash Pipeline

如果你的Jaeger部署已经使用Logstash作为数据中转,或者需要更复杂的处理逻辑,Logstash是很好的选择。

步骤说明

  1. 调整Jaeger输出:将Jaeger的Span输出指向Logstash,而不是直接写入Elasticsearch。
  2. 配置Logstash Pipeline:
    • 输入:接收Jaeger的Span数据(比如用http或kafka输入插件)。
    • 过滤:识别带有phase=end的Span,使用elasticsearch插件查询同traceID的所有Span,获取最小的startTime。
    • 计算:生成endTime和总duration字段。
    • 输出:将原始Span写入Jaeger的ES索引,同时将聚合文档写入报告索引。

简化的Logstash配置示例

input {
  http {
    port => 5044
    codec => json
  }
}

filter {
  # 过滤出phase=end的Span
  if [tags][key] == "phase" and [tags][value] == "end" {
    # 查询当前traceID的所有Span,获取最小startTime
    elasticsearch {
      hosts => ["http://elasticsearch:9200"]
      index => "jaeger-span-*"
      query => '{"query": {"term": {"traceID": "%{[fields][traceID][0]}"}}, "aggs": {"min_start": {"min": {"field": "fields.startTime"}}}, "size": 0}'
      target => "trace_metadata"
    }
    
    # 计算endTime和总duration
    mutate {
      add_field => {
        "[fields][endTime]" => "%{[fields][startTime][0]} + %{[fields][duration][0]}"
        "[fields][total_duration]" => "%{[fields][endTime]} - %{[trace_metadata][aggregations][min_start][value]}"
      }
      convert => {
        "[fields][endTime]" => "integer"
        "[fields][total_duration]" => "integer"
      }
    }
  }
}

output {
  # 写入原始Jaeger索引
  elasticsearch {
    hosts => ["http://elasticsearch:9200"]
    index => "jaeger-span-%{+YYYY-MM-dd}"
    document_id => "%{[@metadata][_id]}"
  }
  
  # 仅将phase=end的聚合文档写入报告索引
  if [tags][key] == "phase" and [tags][value] == "end" {
    elasticsearch {
      hosts => ["http://elasticsearch:9200"]
      index => "jaeger-aggregated-reports"
      document_id => "%{[fields][traceID][0]}"
      document => {
        "fields" => {
          "traceID" => "%{[fields][traceID]}"
          "startTime" => "%{[trace_metadata][aggregations][min_start][value]}"
          "endTime" => "%{[fields][endTime]}"
          "duration" => "%{[fields][total_duration]}"
          "process.serviceName" => "%{[fields][process.serviceName]}"
        }
      }
    }
  }
}

优点

  • 支持更复杂的处理逻辑(比如额外的字段转换、数据清洗)。
  • 不需要依赖Elasticsearch的Watcher功能,适合没有启用X-Pack的环境。

方案3:自定义应用程序(高度定制化)

如果以上方案都无法满足你的需求(比如需要和其他系统集成、复杂的业务规则),可以开发一个自定义的服务来处理。

步骤说明

  1. 轮询ES索引:定期查询jaeger-span-*索引,过滤出phase=end且未处理过的Span(可以用一个单独的索引jaeger-processed-traces来记录已处理的traceID)。
  2. 查询关联Span:对每个未处理的traceID,查询该trace下的所有Span,找到最小的startTime。
  3. 计算并写入:计算endTime和总duration,将聚合文档写入报告索引,同时标记该traceID为已处理。
  4. 处理幂等性:确保同一个traceID不会被重复处理(比如写入报告索引时用traceID作为文档ID,重复写入会自动更新)。

示例Python代码片段(简化版)

from elasticsearch import Elasticsearch
import time

es = Elasticsearch("http://elasticsearch:9200")
REPORT_INDEX = "jaeger-aggregated-reports"
PROCESSED_INDEX = "jaeger-processed-traces"

def process_end_spans():
    # 查询最近5分钟的phase=end的Span,且未处理过
    query = {
        "bool": {
            "must": [
                {"term": {"tags.key": "phase"}},
                {"term": {"tags.value": "end"}}
            ],
            "filter": [
                {"range": {"@timestamp": {"gte": "now-5m"}}},
                {"bool": {"must_not": {"exists": {"field": "processed"}}}}
            ]
        }
    }
    
    resp = es.search(index="jaeger-span-*", query=query, size=100)
    for hit in resp["hits"]["hits"]:
        trace_id = hit["_source"]["fields"]["traceID"][0]
        
        # 查询该trace的最小startTime
        agg_query = {
            "query": {"term": {"traceID": trace_id}},
            "aggs": {"min_start": {"min": {"field": "fields.startTime"}}},
            "size": 0
        }
        agg_resp = es.search(index="jaeger-span-*", body=agg_query)
        min_start = agg_resp["aggregations"]["min_start"]["value"]
        
        # 计算字段
        start_time = hit["_source"]["fields"]["startTime"][0]
        duration = hit["_source"]["fields"]["duration"][0]
        end_time = start_time + duration
        total_duration = end_time - min_start
        
        # 写入报告索引
        es.index(
            index=REPORT_INDEX,
            id=trace_id,
            body={
                "fields": {
                    "traceID": [trace_id],
                    "startTime": [min_start],
                    "endTime": [end_time],
                    "duration": [total_duration],
                    "process.serviceName": hit["_source"]["fields"]["process.serviceName"]
                }
            }
        )
        
        # 标记为已处理
        es.update(index=hit["_index"], id=hit["_id"], body={"doc": {"processed": True}})

if __name__ == "__main__":
    while True:
        process_end_spans()
        time.sleep(300)  # 每5分钟执行一次

优点

  • 完全自定义,适合复杂的业务场景。
  • 可以轻松集成到现有系统中。

方案对比

方案复杂度维护成本灵活性适用场景
Elasticsearch Watcher低低中等简单需求,已有X-Pack权限
Logstash Pipeline中等中等高需要数据中转/清洗,无X-Pack
自定义应用程序高高极高高度定制化需求,系统集成

内容的提问来源于stack exchange,提问作者Vinicius Santos

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 17:07:46