ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Telegraf TopK 处理器深度解析:按聚合函数筛选 Top N 指标序列

Telegraf TopK 处理器深度解析:按聚合函数筛选 Top N 指标序列 Telegraf TopK 处理器深度解析按聚合函数筛选 Top N 指标序列【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegrafTopK 是 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 个字段、k10 会返回 10 个桶而 2 个字段、k3 会返回 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 []参数速查表参数默认值可选值/格式作用period10秒整数两次聚合之间的间隔k10整数每个字段返回的桶数group_by[*]字符串数组支持 glob参与分组聚合的标签空列表表示不按标签聚合fields[value]字符串数组参与聚合的字段aggregationmeansum、mean、min、max聚合函数bottomkfalsetrue/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,pid2088,process_nameXorg cpu_usage7.296576662282613 1546473820000000000 - procstat,pid2780,process_nameibus-engine-simple cpu_usage0 1546473820000000000 - procstat,pid2554,process_namegsd-sound cpu_usage0 1546473820000000000 - procstat,pid3484,process_namechrome cpu_usage4.274300361942799 1546473820000000000 - procstat,pid2467,process_namegnome-shell-calendar-server cpu_usage0 1546473820000000000 - procstat,pid2525,process_namegvfs-goa-volume-monitor cpu_usage0 1546473820000000000 - procstat,pid2888,process_namegnome-terminal-server cpu_usage1.0224991500287577 1546473820000000000 - procstat,pid2454,process_nameibus-x11 cpu_usage0 1546473820000000000 - procstat,pid2564,process_namegsd-xsettings cpu_usage0 1546473820000000000 - procstat,pid12184,process_namedocker cpu_usage0 1546473820000000000 - procstat,pid2432,process_namepulseaudio cpu_usage9.892858669796528 1546473820000000000 --- procstat,pid2432,process_namepulseaudio cpu_usage11.486933087507786 1546474120000000000 procstat,pid2432,process_namepulseaudio cpu_usage10.056503212060552 1546474130000000000 procstat,pid23620,process_namechrome cpu_usage2.098690278123081 1546474120000000000 procstat,pid23620,process_namechrome cpu_usage17.52514619948493 1546474130000000000 procstat,pid2088,process_nameXorg cpu_usage1.6016732172309973 1546474120000000000 procstat,pid2088,process_nameXorg cpu_usage8.481040931533833 1546474130000000000可以看到原始输入中大量低 CPU 占用cpu_usage0的进程被过滤掉只有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,hostweb-01 cpu_usage42.5,cpu_usage_topk_rank1,cpu_usage_topk_aggregate41.8,topk_groupcpuhostweb-01 1690000000000000000GroupBy 标签的格式在源码中有明确定义由测量名、分隔符与tagvalue键值对拼接而成且标签键经过排序以保证键的确定性见 generateGroupByKey。这一点也被 TestTopkGroupByKeyTag 直接验证例如期望值metric1tag1TWOtag3SIX。源码级原理从分组键到排序输出分组键的生成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 语法是可行的。聚合函数的实现getAggregationFunctiontopk.go为四种聚合分别生成闭包sum对组内各指标的字段值直接累加min / max分别以math.MaxFloat64/-math.MaxFloat64为初值做极值比较mean先求和并计数最后除以样本数若某个字段在整个周期内没有任何样本则聚合值记为 0。数值转换由convert函数完成topk.go仅支持float64、int64、uint64三种类型遇到无法转换的字段值会通过日志记录一条Cannot convert value ...信息并跳过该值不影响整体聚合。排序与去重sortMetricstopk.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、debugmetric 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.044sum 为140.22TestTopkGroupby1/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),仅供参考
返回列表