如何在Elasticsearch中基于Jaeger的插入文档生成聚合新文档?
解决方案:实现Jaeger Span的聚合报告文档生成
你提到想用CDC但Elasticsearch不支持原生CDC,这里有几个可行的方案来完成你的需求,从简单到复杂依次介绍:
方案1:使用Elasticsearch Watcher(原生告警功能)
这是最轻量化的方案,利用Elasticsearch自带的Watcher功能来监听新插入的phase=end文档,并自动完成聚合逻辑。
步骤说明
- 创建Watcher:配置一个Watcher,监听
jaeger-span-*索引的文档新增事件,过滤出带有tags.key:phase且tags.value:end的文档。 - 查询关联的初始Span:在Watcher的执行阶段,通过Search查询同一个
traceID下的所有Span,找到最小的startTime值(也就是初始文档的startTime)。 - 计算聚合字段:
- 计算
endTime:当前phase=end文档的startTime+duration - 计算总
duration:endTime- 初始Span的startTime
- 计算
- 写入聚合文档:使用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是很好的选择。
步骤说明
- 调整Jaeger输出:将Jaeger的Span输出指向Logstash,而不是直接写入Elasticsearch。
- 配置Logstash Pipeline:
- 输入:接收Jaeger的Span数据(比如用
http或kafka输入插件)。 - 过滤:识别带有
phase=end的Span,使用elasticsearch插件查询同traceID的所有Span,获取最小的startTime。 - 计算:生成
endTime和总duration字段。 - 输出:将原始Span写入Jaeger的ES索引,同时将聚合文档写入报告索引。
- 输入:接收Jaeger的Span数据(比如用
简化的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:自定义应用程序(高度定制化)
如果以上方案都无法满足你的需求(比如需要和其他系统集成、复杂的业务规则),可以开发一个自定义的服务来处理。
步骤说明
- 轮询ES索引:定期查询
jaeger-span-*索引,过滤出phase=end且未处理过的Span(可以用一个单独的索引jaeger-processed-traces来记录已处理的traceID)。 - 查询关联Span:对每个未处理的traceID,查询该trace下的所有Span,找到最小的
startTime。 - 计算并写入:计算
endTime和总duration,将聚合文档写入报告索引,同时标记该traceID为已处理。 - 处理幂等性:确保同一个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
相关产品推荐
相关产品推荐

