Telegraf cumulative_sum 处理器插件详解:按序列对字段做累积求和
【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf
本指南围绕 Telegraf 仓库中的cumulative_sum处理器插件(源码见 cumulative_sum.go,官方文档见 README.md)展开,完整讲解它的作用、配置方式、运行机制与使用边界。该插件会在每个 metric 被更新时为指定数值字段生成持续递增的累积和(xxx_sum),非常适合为依赖“单调递增数值”的输出端(如计算速率的监控后端)提供前置数据处理。读完本文,你将掌握cumulative_sum的字段筛选、缓存过期策略、通配符用法,以及它与普通字段累加/聚合器之间的本质区别。
插件元信息:Telegraf v1.35.0引入 · 类别transformation(转换)· 支持平台all(全平台)。
插件概述与适用场景
cumulative_sum是一个典型的状态型转换处理器:它按 metric 序列(series)在内存中缓存历史字段值,每当收到一条新 metric 时,就把该字段的历史累计值加上当前值,并以新字段<原字段名>_sum追加到 metric 上输出。
它的典型应用场景是配合需要单调递增数值的输出端。例如某些输出系统要求指标必须是持续递增的计数器(counter),而你的采集源提供的是“瞬时增量”或“周期快照”,此时就可以用cumulative_sum把增量逐步累加,把普通数值伪装成单调递增的序列,从而满足输出端的数据模型要求。
[!NOTE] 同一序列(series)内的 metric 是按到达顺序(order of arrival)累加的,而不是按时间戳(timestamp)排序后累加!这意味着乱序到达的指标会按实际到达次序累加,累加结果与指标自身时间戳的顺序无关。这一点在高并发、乱序传递的管道场景中需要特别注意。
配置详解
cumulative_sum的完整配置如下(与仓库中的 sample.conf 完全一致):
# Compute the cumulative sum of the given fields [[processors.cumulative_sum]] ## Numerical fields to be processed (accepting wildcards) # fields = ["*"] ## Interval after which metrics are evicted from the cache and the ## sum values are reset to zero. A zero or unset value will keep the ## metric forever. ## It is strongly recommended to set an expiry interval to avoid ## growing memory usage when varying metric series are processed. # expiry_interval = "0s"参数一:fields —— 参与累加的字段列表
- 类型:字符串数组(
[]string),对应源码结构体中的Fields []string \toml:"fields"`` 字段。 - 默认值:
["*"],即对 metric 中所有字段都尝试累加。 - 通配符:支持 glob 风格的通配符,例如
fields = ["bytes_*"]可只处理以bytes_开头的字段。 - 实现细节:在源码的
Init()方法中,若Fields为空会默认赋值为["*"],随后通过filter.Compile(c.Fields)(来自 filter 包)编译成匹配器,供每次处理时用c.accept.Match(field.Key)判断某字段是否参与累加。
参数二:expiry_interval —— 缓存过期时间
- 类型:duration 字符串(如
"10m"、"1h"),对应ExpiryInterval config.Duration。 - 默认值:
0s(未设置)。此时缓存条目永不失效,累加值会一直累计下去。 - 作用:设置后,超过该时长未被再次观测到的序列缓存条目会被从内存中清除,其累计和归零;当该序列后续再次出现时,将从新基数重新开始累加。
- 强烈建议配置:当采集的 metric 序列组合不断变化(例如 tag 值持续新增)时,若不设置过期时间,缓存会随序列数量增长而持续占用内存。官方文档与源码注释均明确建议设置合理的过期区间以避免内存无限增长。
一个面向实战的推荐配置示例:
[[processors.cumulative_sum]] fields = ["bytes_sent", "bytes_received", "packets_*"] expiry_interval = "10m"工作原理:从源码看累加流程
核心逻辑集中在Apply()方法中(cumulative_sum.go),整个处理流程如下:
- 定位序列缓存:对每条输入 metric 调用
original.HashID()计算序列标识(由测量名 + tag 集合决定的哈希值),在cache map[uint64]*entry中查找该序列的历史累计值;若不存在则新建一个空的entry。 - 复制 metric:通过
original.Copy()生成副本,保证原始 metric 不被修改,同时保留原有全部字段与 tag。 - 字段筛选与转换:遍历副本的
FieldList():- 不匹配
fields过滤器(或未配置)的字段直接保留、不处理; - 对匹配字段调用
internal.ToFloat64()(实现在 internal/type_conversions.go)尝试转为float64。该函数支持string、[]byte、bool、各种有符号/无符号整数、浮点数等类型(如bool会转为 1/0),转换失败(如结构体、切片等不支持类型)的字段会被跳过,并输出一条 Trace 级别日志; - 计算
sum := stored.sums[field.Key] + fv,用m.AddField(field.Key+"_sum", sum)追加<字段名>_sum字段,并回写缓存。
- 不匹配
- 更新与返回:刷新该条目的
seen时间戳,把更新后的 entry 写回 cache,输出携带_sum字段的新 metric,并调用original.Accept()标记原 metric 已被消费。 - 过期清理:若
ExpiryInterval > 0,以当前时间为基准计算阈值,通过maps.DeleteFunc删除所有seen早于阈值的缓存条目。
插件通过包内init()函数调用processors.Add("cumulative_sum", ...)完成注册,并经由 plugins/processors/all/cumulative_sum.go 统一加载(构建标签!custom || processors || processors.cumulative_sum),因此它已包含在标准发行版中,无需额外编译。
官方示例解析
README 中给出的示例直观展示了处理前后的差异(其中-为输入,+为输出):
- net,host=server01 bytes_sent=1000,bytes_received=500 - net,host=server01 bytes_sent=2500,bytes_received=1500 - net,host=server01 bytes_sent=3000,bytes_received=2500 + net,host=server01 bytes_sent=1000,bytes_sent_sum=1000,bytes_received=500,bytes_received_sum=500 + net,host=server01 bytes_sent=2500,bytes_sent_sum=3500,bytes_received=1500,bytes_received_sum=2000 + net,host=server01 bytes_sent=3000,bytes_sent_sum=6500,bytes_received=2500,bytes_received_sum=4500从中可以看到三个关键行为:
- 保留原始字段:
bytes_sent、bytes_received原样保留,_sum只是追加字段,不会覆盖原始值; - 逐条累积:同一序列(
net,host=server01)的第二条输入中,bytes_sent_sum = 1000 + 2500 = 3500,第三条为3500 + 3000 = 6500,呈现单调递增趋势; - 按到达顺序累加:示例输入本身按时间先后到达,因此累加顺序与时间顺序一致;但若 metric 乱序到达,累加依然以到达顺序为准(见上文 NOTE)。
测试用例佐证:行为边界一目了然
仓库中的 cumulative_sum_test.go 用三组用例精确定义了插件的行为边界,是对文档的有力补充:
- all fields keep original fields:默认
fields为 nil 时所有字段都会被处理。测试显示,bool字段healthy=false被转换为healthy_sum=0(对应ToFloat64的 bool 处理规则);int64字段error_counter=10得到error_counter_sum=10;字符串字段error="machine broken"无法转浮点,被跳过且不生成_sum。这印证了“字段可转为 float 才会被累加”的实现事实。 - filter value remove original:配置
fields=["value"]后,只有value字段生成value_sum,其余字段(如healthy、error_counter、error)原样保留、不参与累加。 - multiple metrics:同一测量名 + 不同 tag 的序列各自独立累加(
tag=some tag累计到 3.0,tag=another tag保持 4.4),证明缓存是以HashID()(测量名 + tag 组合)为粒度的。 - TestCacheExpiry:模拟
expiry_interval=10s时,若某序列超过 10 秒未被观测,其缓存条目被清除、累计值归零,后续再出现时从新基数累加(4.4不再叠加进3.2,而是重新从4.4开始)。
使用注意与边界
综合文档与源码,使用cumulative_sum时需要注意以下几点:
- 乱序问题:按到达顺序累加而非按时间戳累加。若你的数据源存在乱序投递且对累加顺序敏感,请先在上游保证有序,或在管道中用其他机制排序。
- 非数值字段:只有能通过
internal.ToFloat64转成float64的字段才会被累加(bool会被转成 0/1)。无法转换的字段被静默跳过(仅记录 Trace 日志),不会报错中断。 - 内存管理:强烈建议设置
expiry_interval。否则当 tag 组合持续变化时,缓存条目只增不减,内存占用会随序列规模线性增长。设置过期时间后,长期未出现的序列会被回收,其累计和自动归零。 - 与聚合器(aggregator)的区别:
cumulative_sum是流式转换处理器,每收到一条 metric 就立即输出携带累计值的 metric,不依赖时间窗口;而聚合器(如minmax、basicstats等)通常在窗口结束后才批量输出。它也不等于“对字段求和后再输出单一值”——它保留原始字段并逐条追加_sum。 - 处理器顺序:
cumulative_sum属于 processors 阶段,其执行顺序受全局插件顺序配置影响,具体可参考 docs/CONFIGURATION.md 中的插件排序说明;该处理器同样支持所有全局插件配置项(如namepass、tagexclude等,详见 docs/includes/plugin_config.md)。
总结
cumulative_sum是 Telegraf 中实现“按序列单调递增累加”最直接的处理器:一条[[processors.cumulative_sum]]配置,配合fields通配筛选与expiry_interval内存治理,即可把瞬时增量转换为单调递增计数,适配依赖 counter 语义的输出端。理解其“按到达顺序累加、按序列独立缓存、按过期时间回收”三大机制,你就能在真实管道中安全、可控地使用它。
【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考