Telegraf TopK 处理器深度解析:按聚合函数筛选 Top N 指标序列
【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf
TopK 是 Telegraf 中一个典型的"变换(transformation)"型处理器插件(自 v1.7.0 起提供),用于在一段周期内对指标按测量名与标签进行分组,对指定字段施加sum/mean/min/max聚合后,只放行排名最靠前的K个分组(bucket)。本文基于当前仓库中的 README.md 展开,结合 topk.go 源码与 topk_test.go 测试用例,完整讲解其全部配置参数、运行机制、实战配置与底层实现细节,帮助你用它实现"只看最热的 K 条时序"这类降噪与聚焦需求。
功能定位:它解决什么问题
在监控场景中,输入源(如procstat、exec、snmp等)常常产生成百上千条时序。你真正关心的可能只是其中"最突出"的一部分——例如 CPU 占用最高的几个进程、流量最大的几个网络接口。TopK 处理器就是为此设计的:它不会丢弃所有数据,而是周期性聚合并仅输出聚合结果排名前 K 的指标组,同时支持bottomk反转为保留最低的 K 组。
从源码结构看,它内部维护一个cache map[string][]telegraf.Metric(见 topk.go),每个分组键对应一批待聚合的指标;每个period周期结束时统一计算并一次性输出,行为上更接近"按周期刷新的聚合型处理器"。
处理流程:分组 → 聚合 → 取 Top K
文档明确给出了处理器对一批指标执行的三步流程(见 README.md):
- 分组:依据指标的测量名(metric name)与标签(tags),把指标归入对应的桶(bucket)。分组键由测量名与参与分组的标签拼接而成。
- 聚合:每
period秒,对每个桶、每个被选中的字段,使用指定的聚合函数(min、sum、mean、max)计算聚合值。 - 取 Top K:对每个字段的聚合结果按值排序,放行排名前
K的桶内的全部指标。
注意步骤 3 的关键语义:输出的是"排名前 K 的桶里的所有原始指标",而不是 K 条聚合后的指标。因此当某个桶内指标数量较多时,实际输出的系列数可能多于 K(详见"注意事项"一节)。
在 topk.go 的Apply中可以看到完整实现:
- 每批到达的指标先检查是否包含
fields中声明的任一字段(m.HasField(f)),一个都不包含则直接丢弃; - 通过
groupBy写入内部缓存,并判断距上次聚合的时间time.Since(t.lastAggregation) >= period是否到期; - 到期则调用
push()完成排序输出,未到期则返回nil(指标被缓存等待下一个周期)。
完整配置与参数详解
处理器完整示例配置见 sample.conf,与 README 中展示的配置一致。以下是带注释的完整配置:
[[processors.topk]] ## 每次聚合之间的时间间隔(秒) # period = 10 ## 每个字段返回的 top 桶数量 ## 每一个声明参与聚合的字段都会独立返回 k 个结果。 ## 例如:1 个字段、k=10 会返回 10 个桶;而 2 个字段、k=3 会返回 6 个桶。 # k = 10 ## 聚合所依据的标签。支持 glob 通配符,匹配到的任意标签都会参与聚合。 ## 若设置为空列表,则完全不按标签聚合。 # group_by = ['*'] ## 参与聚合的字段 ## 每个字段都会生成一个独立的聚合,每次聚合返回 k 个桶。 ## 若某条指标不包含该字段,则该指标会从这次聚合中被丢弃。 ## 如有需要,可考虑配合 defaults 处理器插件预先补齐字段。 # fields = ["value"] ## 使用的聚合函数。可选值:sum、mean、min、max # aggregation = "mean" ## 若为 true,则返回聚合值最低的 k 个桶(bottom k),而不是最高的 k 个 # bottomk = false ## 插件会为每条指标生成一个由其测量名与标签计算出的 GroupBy 标签。 ## 若该设置非空字符串,插件会以该设置的值作为标签名, ## 把计算出的 GroupBy 值作为标签值附加到每条指标上,便于调试。 # add_groupby_tag = "" ## 用于获取每条指标在 top k 中排名位置的字段设置。 ## 'add_rank_fields' 指定需要排名的字段;若列表非空,则对列表中每个字段, ## 向每条指标添加一个字段,其值为该指标所属分组在该字段聚合结果中的排名。 ## 字段名 = 聚合字段名 + 后缀 '_topk_rank' # add_rank_fields = [] ## 用于获取聚合值的字段设置。 ## 'add_aggregate_fields' 指定需要聚合值的字段;若列表非空,则对列表中每个字段, ## 向每条指标添加一个字段,其值为该指标所属分组的最终聚合结果。 ## 字段名 = 聚合字段名 + 后缀 '_topk_aggregate' # add_aggregate_fields = []参数速查表
| 参数 | 默认值 | 可选值/格式 | 作用 |
|---|---|---|---|
period | 10(秒) | 整数 | 两次聚合之间的间隔 |
k | 10 | 整数 | 每个字段返回的桶数 |
group_by | ['*'] | 字符串数组,支持 glob | 参与分组聚合的标签;空列表表示不按标签聚合 |
fields | ["value"] | 字符串数组 | 参与聚合的字段 |
aggregation | mean | sum、mean、min、max | 聚合函数 |
bottomk | false | true/false | 是否返回最低的 k 个桶 |
add_groupby_tag | "" | 字符串 | 非空时给每条指标附加 GroupBy 标签 |
add_rank_fields | [] | 字符串数组 | 非空时为每条指标附加_topk_rank排名字段 |
add_aggregate_fields | [] | 字符串数组 | 非空时为每条指标附加_topk_aggregate聚合值字段 |
默认值在 topk.go 的newTopK()中直接可见:period = 10s、k = 10、fields = ["value"]、aggregation = "mean"、group_by = ["*"],并调用Reset()初始化缓存与计时起点。
关于add_rank_fields与add_aggregate_fields的命名规则
- 排名字段:聚合字段名 +
_topk_rank。例如对字段cpu_usage聚合,则每条被输出指标会获得字段cpu_usage_topk_rank,值为该分组在本次排名中的名次(从 1 开始)。 - 聚合值字段:聚合字段名 +
_topk_aggregate。例如cpu_usage_topk_aggregate,值为该分组的聚合计算结果。
这两个机制在 push() 中实现:只有当addRankFields/addAggregateFields非空时才逐条为输出指标附加字段,且要求该指标确实拥有对应聚合字段(m.HasField(field))才会写入。
实战示例:找出 CPU 占用最高的进程
README 提供了一个非常直观的例子:用procstat采集各进程的cpu_usage,只关心占用最高的 3 个进程:
[[processors.topk]] period = 20 k = 3 group_by = ["pid"] fields = ["cpu_usage"]这里:
period = 20:每 20 秒做一次聚合排名;k = 3:每个字段返回 3 个桶;group_by = ["pid"]:按pid标签分组,即同一进程不同采集点的指标被归入同一组;fields = ["cpu_usage"]:仅对cpu_usage字段做聚合(默认mean)。
输出前后对比
README 中展示了启用该处理器前后的数据差异(-为原始输入,+为输出,时间戳已截取):
- procstat,pid=2088,process_name=Xorg cpu_usage=7.296576662282613 1546473820000000000 - procstat,pid=2780,process_name=ibus-engine-simple cpu_usage=0 1546473820000000000 - procstat,pid=2554,process_name=gsd-sound cpu_usage=0 1546473820000000000 - procstat,pid=3484,process_name=chrome cpu_usage=4.274300361942799 1546473820000000000 - procstat,pid=2467,process_name=gnome-shell-calendar-server cpu_usage=0 1546473820000000000 - procstat,pid=2525,process_name=gvfs-goa-volume-monitor cpu_usage=0 1546473820000000000 - procstat,pid=2888,process_name=gnome-terminal-server cpu_usage=1.0224991500287577 1546473820000000000 - procstat,pid=2454,process_name=ibus-x11 cpu_usage=0 1546473820000000000 - procstat,pid=2564,process_name=gsd-xsettings cpu_usage=0 1546473820000000000 - procstat,pid=12184,process_name=docker cpu_usage=0 1546473820000000000 - procstat,pid=2432,process_name=pulseaudio cpu_usage=9.892858669796528 1546473820000000000 --- + procstat,pid=2432,process_name=pulseaudio cpu_usage=11.486933087507786 1546474120000000000 + procstat,pid=2432,process_name=pulseaudio cpu_usage=10.056503212060552 1546474130000000000 + procstat,pid=23620,process_name=chrome cpu_usage=2.098690278123081 1546474120000000000 + procstat,pid=23620,process_name=chrome cpu_usage=17.52514619948493 1546474130000000000 + procstat,pid=2088,process_name=Xorg cpu_usage=1.6016732172309973 1546474120000000000 + procstat,pid=2088,process_name=Xorg cpu_usage=8.481040931533833 1546474130000000000可以看到,原始输入中大量低 CPU 占用(cpu_usage=0)的进程被过滤掉,只有pulseaudio、chrome、Xorg三个进程(按pid分组后排名前三的组)内的全部指标被保留,且每条输出的时间戳来自采集时刻,而不是聚合时刻——输出仍保留原始采集点,只是数据量被大幅收敛。
进阶用法:排名与聚合值附加
在实际告警或可视化场景中,你可能希望知道"这条指标属于第几名"或"它所在分组算出的聚合值是多少"。此时组合使用add_rank_fields与add_aggregate_fields:
[[processors.topk]] period = 30 k = 5 group_by = ["host"] fields = ["cpu_usage"] aggregation = "mean" add_rank_fields = ["cpu_usage"] # 输出 cpu_usage_topk_rank add_aggregate_fields = ["cpu_usage"] # 输出 cpu_usage_topk_aggregate add_groupby_tag = "topk_group" # 输出 topk_group 标签,值为 GroupBy 键配置后每条输出的指标大致形如:
cpu,host=web-01 cpu_usage=42.5,cpu_usage_topk_rank=1,cpu_usage_topk_aggregate=41.8,topk_group="cpu&host=web-01&" 1690000000000000000GroupBy 标签的格式在源码中有明确定义:由测量名、&分隔符与tag=value&键值对拼接而成,且标签键经过排序以保证键的确定性(见 generateGroupByKey)。这一点也被 TestTopkGroupByKeyTag 直接验证,例如期望值"metric1&tag1=TWO&tag3=SIX&"。
源码级原理:从分组键到排序输出
分组键的生成
generateGroupByKey使用 filter/ 包的filter.Compile把group_by中的 glob 表达式编译为匹配器(支持*通配符),再遍历指标标签,仅保留匹配的标签拼接为分组键。因此:
group_by = ['*']:所有标签都参与分组;group_by = ["pid"]:仅pid标签参与;group_by = []:完全不按标签分组,只按测量名分组——TestTopkGroupbyMetricName1 验证了只按测量名分组的行为;- 还可以使用带通配符的表达式,如测试中的
"tag[13]"、"tag[12]"(topk_test.go),说明tag[13]这类 glob 语法是可行的。
聚合函数的实现
getAggregationFunction(topk.go)为四种聚合分别生成闭包:
- sum:对组内各指标的字段值直接累加;
- min / max:分别以
math.MaxFloat64/-math.MaxFloat64为初值做极值比较; - mean:先求和并计数,最后除以样本数;若某个字段在整个周期内没有任何样本,则聚合值记为 0。
数值转换由convert函数完成(topk.go),仅支持float64、int64、uint64三种类型;遇到无法转换的字段值,会通过日志记录一条Cannot convert value ...信息并跳过该值,不影响整体聚合。
排序与去重
sortMetrics(topk.go)使用sort.SliceStable做稳定排序,保证并列时顺序确定;bottomk = true时按升序取前 K(最低的 K 组),否则按降序取前 K。随后push()通过addedKeys记录已加入输出结果的分组键,确保同一个组在多个字段的 Top K 中不会重复输出。
此外,处理器对输出的指标会调用metric.New重建新指标对象(topk.go),等价于对输入做去重处理——这也是 README 中"处理器会对指标去重"这一说明的实现依据。
与 tracking metric 的配合
在 Apply 的入口处,每条输入指标都会被调用m.Accept()。注释解释了原因:若缓存持有了带跟踪(tracking)的指标而不及时确认投递,可能阻塞输入端等待回执。因此处理器把所有收到的指标视为"已投递",输出时再以未跟踪的新指标形式向下游传递。TestTracking 验证了这一行为:在单个周期内输出指标数与输入一致,且所有原始 tracking metric 都能收到投递回执。
注意事项与边界情况
综合 README 的 Notes 与源码行为,使用时有以下几点需要牢记:
- 输出数量可能超过 K:处理器按"桶"排名,每个桶内的所有原始指标都会被放行。桶内指标多时,实际输出的系列数会多于 K。
- 缺字段的指标会被丢弃:若一条指标不包含
fields中声明的任何一个字段,它会被直接排除在该次聚合之外。README 建议:如需保证字段存在,可先在管道中放置 defaults 处理器插件补齐字段。 - 测量名始终参与分组:即使
group_by为空列表,分组键仍会包含测量名,因此不同测量名的指标永远不会混入同一组。 - 指标默认不被修改:处理器默认不添加任何标签与字段(见 README 的 Tags/Fields 小节),仅当
add_groupby_tag、add_rank_fields、add_aggregate_fields被设置为非空值时才附加相应信息。 - 时间窗口语义:缓存内的指标会一直累积,直到距上次聚合超过
period才统一结算并清空缓存(Reset())。这意味着输出会存在一个周期级别的延迟,适合周期性汇总场景,不适合实时逐条转发。 - 组内聚合缺失字段的处理:单个周期内某字段无任何样本时,mean 聚合结果记为 0,实际测试与文档一致。
在管道中的位置与通用配置
TopK 属于 processors 插件类别,默认在输入插件之后、聚合器插件之前执行(参见 docs/CONFIGURATION.md)。文档的 Global configuration options 指出,所有处理器都支持通用全局配置项,例如:
- alias:为插件实例命名;
- order:多处理器时的执行顺序(从 1 开始),未指定时按配置文件中出现的顺序执行;
- log_level:覆盖该插件的日志级别(
error、warn、info、debug); - metric filtering 参数:用于限定哪些指标进入该处理器。
若管道中同时存在多个处理器且顺序敏感,须对所有相关处理器显式设置order,示例见 docs/CONFIGURATION.md。更详尽的处理器通用说明可参考 docs/PROCESSORS.md。此外,docs/AGGREGATORS_AND_PROCESSORS.md 对处理器与聚合器的差异(聚合器产出新指标、处理器改写原指标)也有系统阐述,可作为理解 TopK 定位的补充阅读。
验证与测试
仓库为 TopK 提供了覆盖较全的单元测试(topk_test.go),可作为理解行为与二次开发的参考:
TestTopkAggregatorsSmokeTests:四种聚合函数的冒烟测试;TestTopkMeanAddAggregateFields/Sum/Max/Min:分别验证四类聚合下_topk_aggregate字段的取值(例如 mean 组内 5 条指标聚合值为28.044,sum 为140.22);TestTopkGroupby1/2/3、TestTopkGroupbyFields1/2:验证分组键、glob 标签(如tag[13])、多字段独立聚合;TestTopkGroupbyMetricName1/2:验证按测量名分组;TestTopkBottomk:验证bottomk反选最低 K 组;TestTopkGroupByKeyTag:验证add_groupby_tag附加的 GroupBy 键值格式;TestTracking:验证 tracking metric 的投递确认行为。
小结
TopK 处理器以"周期聚合 + Top K 排序输出"的方式,为海量时序提供了一种轻量、可配置的聚焦手段。核心要点可概括为:用group_by决定分组维度、用fields决定聚合字段、用aggregation决定排序依据、用k与bottomk决定保留规模,用add_rank_fields/add_aggregate_fields/add_groupby_tag为输出附加排名、聚合值与分组信息。其实现细节(分组键生成、稳定排序、去重输出、tracking 指标处理)都在 topk.go 中清晰可查,测试用例则提供了行为层面的完整背书,是理解 Telegraf 处理器数据流与变换型插件设计的上佳范本。
【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考