Telegraf 插件内部统计(internal plugin statistics):基于 selfstat.Collector 的插件级指标采集规范与实践
【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf
本篇文章围绕 Telegraf 仓库中的设计规范文档 tsd-011-internal-plugin-statistics.md 展开,系统讲解如何让插件上报"只有插件实例自身才知道"的内部统计指标(例如输出插件实际写入的字节数、特定错误类型等),并通过internal输入插件对外暴露。文中将结合selfstat包源码、models层的自动注入机制以及真实插件(如influxdb输出插件)的落地实现,说明从"定义一个Statistics字段"到"指标进入常规采集管线"的完整链路,帮助读者掌握编写或扩展带内部统计能力的 Telegraf 插件的方法。
一、为什么需要插件级内部统计
Telegraf 通过 internal 输入插件 提供对运行中插件的统计能力,这些指标是运维与调优 Telegraf 部署、定位问题的重要依据。从源码实现看,该插件定义了三类采集内容(见 plugins/inputs/internal/internal.go):
internal_memstats:Go 运行时内存统计(默认开启,collect_memstats = true);internal_gostats:Goruntime/metrics暴露的指标(默认关闭,collect_gostats = false);- 插件统计:通过
selfstat.Metrics()汇聚所有已注册的自统计指标,默认按插件类型聚合输出,也可通过per_instance = true改为逐实例输出。
其中"插件统计"又分为两类:
- 模型层(model-level)统计:由
models包中的各 Running 类(如running_input.go、running_output.go)负责注册,覆盖采集时长、写入时长、错误次数等通用维度; - 插件内部(plugin-internal)统计:只有插件实例自身知道的信息,例如向输出目标实际写入的字节数、特定协议错误、请求计数等,无法在模型层获知。
问题在于:插件若绕过模型层直接调用selfstat.Register注册统计对象,将无法获得全局tags配置与alias别名设置——这些信息只在模型层掌握。tsd-011 规范的初衷正是解决这一缺陷:定义一个框架,把携带了alias与 tags 信息的"统计收集器"(statistics collector)注入插件,让插件通过它注册内部统计。相关背景可追溯至仓库历史中的 issue #4889(InfluxDB 输出插件内部统计)、#6965(Kafka 输出插件内部统计)与 #17275(InfluxDB v2 输出插件内部统计)。
二、总体机制:注入一个统计收集器
规范提出的核心设计包含三个环节:
- 插件侧:插件结构体导出
Statistics成员,类型必须为指针*selfstat.Collector,强烈建议为其声明"-"的 TOML tag,避免用户配置与内部成员冲突; - 模型层注入:Telegraf 模型代码在实例化插件之后、调用插件
Init函数之前,将selfstat.Collector实例注入Statistics成员,注入时携带模型层已知的alias、tags等全部信息; - 插件使用:插件必须通过该 collector 作为代理来完成统计的注册(register)、注销(unregister)、重置(reset)与访问(access),而不是直接调用全局的
selfstat.Register。
2.1 注入的源码实现
注入逻辑位于 models/common.go 的SetStatisticsOnPlugin函数:
func SetStatisticsOnPlugin(plugin interface{}, logger telegraf.Logger, tags map[string]string) { // Find the statistics collector instance := reflect.Indirect(reflect.ValueOf(plugin)) field := instance.FieldByName("Statistics") if !field.IsValid() { return } // Validate the type and make sure we can actually set the struct field if field.Type().String() != "*selfstat.Collector" || !field.CanSet() { logger.Debugf( "Plugin %q defines a 'Statistics' field on its struct of an unexpected type %q. Expected *selfstat.Collector", instance.Type().Name(), field.Type().String(), ) return } // Create a new collector and set it collector := selfstat.NewCollector(tags) field.Set(reflect.ValueOf(collector)) }该函数使用反射定位插件结构体上的Statistics字段:
- 字段不存在(
!field.IsValid())时静默返回,保证未实现内部统计的插件不受影响; - 字段类型不是
*selfstat.Collector或不可写时,仅输出 Debug 日志并跳过,避免破坏插件启动; - 校验通过后,以模型层构造的
tags创建新的 collector 并写入字段。
这里传入的tags已经包含了模型层信息。以输出插件为例,models/running_output.go 中构造的 tags 为:
tags := map[string]string{ "output": config.Name, "_id": config.ID, } if config.Alias != "" { tags["alias"] = config.Alias } errorLogRegister := selfstat.Register("write", "errors", tags) logger := logging.New("outputs", config.Name, config.Alias) logger.RegisterErrorCallback(func() { errorLogRegister.Incr(1) }) if err := logger.SetLogLevel(config.LogLevel); err != nil { logger.Error(err) } SetLoggerOnPlugin(output, logger) SetStatisticsOnPlugin(output, logger, tags)可见output=<插件名>、_id=<插件实例 ID>、可选的alias=<别名>在模型层被统一封装进 collector,随后通过SetStatisticsOnPlugin注入插件。
2.2 Collector 的职责与能力
selfstat.Collector定义在 selfstat/collector.go:
type Collector struct { tags map[string]string statistics map[string]Stat }它持有一份"收集器级"的 tags(即模型层注入的 alias/tags),并缓存本收集器注册过的统计对象。其对外方法完整对应规范要求的四种能力:
Register(measurement, field string, tags map[string]string) Stat(collector.go):注册普通统计,内部先把传入 tags 与收集器自身 tags合并,再调用全局selfstat.Register,并缓存 key 避免重复注册;RegisterTiming(measurement, field string, tags map[string]string) Stat(collector.go):注册耗时类统计,语义与Register相同,但底层使用RegisterTiming(见下文"普通统计与计时统计");Unregister(measurement, field string, tags map[string]string)(collector.go):注销指定统计并从缓存删除,另有UnregisterAll()批量注销全部统计;Get(measurement, field string, tags map[string]string) Stat与Reset(measurement, field string, tags map[string]string)(collector.go):按 key 读取统计对象,或将指定统计的值重置为 0。
由此,插件侧拿到的是"已预先绑定了模型层 tags 的代理",注册的每个统计都会自动带上alias、插件名等模型层标签,从根上解决了直接调用全局selfstat.Register丢失 tags 的问题。
2.3 全局注册表:collector 的底层支撑
collector 最终仍委托给selfstat包的全局注册表完成实际登记。selfstat/selfstat.go 暴露了核心 API:
Register(measurement, field string, tags map[string]string) Stat:注册普通统计,测量名会被自动加上internal_前缀;RegisterTiming(measurement, field string, tags map[string]string) Stat:注册计时统计,同样带internal_前缀;Unregister(...):从注册表移除统计;Metrics() []telegraf.Metric:把注册表中所有统计转换为 Telegraf 指标,供internal插件采集。
注册表内部以map[uint64]map[string]Stat组织(selfstat.go),外层 key 由"测量名 + 排序后的 tags"经 FNV-1a 哈希生成(key函数,selfstat.go),内层再按字段名索引;并发安全由sync.Mutex保证。Stat接口(selfstat.go)定义了指标对象的核心操作:
Name()/FieldName()/Tags():返回测量名、字段名与 tags(每次调用返回新 map);Incr(v int64):普通统计累加;计时统计则将本次耗时加入缓存;Set(v int64):普通统计直接赋值;计时统计同样写入缓存;Get():读取当前值;计时统计返回"自上次Get()以来所有计时的平均值",无新计时则沿用上次值;Unregister():从注册表移除。
一个典型输出后立刻见到的现象是:插件统计在internal_<plugin_name>测量下按 measurement 聚合多个字段,配合internal插件默认的per_instance = false(按插件类型聚合)或per_instance = true(按实例输出,保留_idtag)两种模式对外呈现。
三、插件侧落地步骤与完整示例
按规范落地一个带内部统计的插件,需要三步:声明字段 → 在Init中注册统计 → 在运行路径上更新统计。
3.1 声明Statistics成员
在插件结构体中增加指针类型的Statistics *selfstat.Collector字段,并打上toml:"-"标签:
type InfluxDB struct { // ... 既有配置字段 ... Statistics *selfstat.Collector `toml:"-"` // 插件自己的统计句柄 bytesWritten selfstat.Stat }toml:"-"的作用是防止用户在配置文件中通过同名键覆盖该成员。这正是规范中"强烈建议定义"-"TOML tag"的落地形态,可在 plugins/outputs/influxdb/influxdb.go 看到真实样例。
3.2 在Init中通过 collector 注册统计
Init是注入完成的保证点:模型层在实例化之后、调用Init之前注入 collector,因此Init内可以安全使用i.Statistics。参考 plugins/outputs/influxdb/influxdb.go:
func (i *InfluxDB) Init() error { // ... 既有默认值与序列化器初始化 ... // Register internal metrics i.bytesWritten = i.Statistics.Register("write", "bytes_written", nil) return nil }这里以 measurementwrite、字段bytes_written注册统计(collector 会拼接internal_前缀,最终测量名为internal_write),tags 传nil表示只使用收集器自身携带的模型层 tags。
3.3 在写入路径上更新统计
注册得到的selfstat.Stat句柄保存在插件字段中,之后在每次成功写入后更新,例如httpClient/udpClient会把BytesWritten: i.bytesWritten传入客户端,由客户端在写出字节后调用Incr(n)(见 plugins/outputs/influxdb/influxdb.go)。这样internal_write测量中就会实时反映该输出实例累计写入的字节数,且天然带上了output、_id、alias等模型层标签。
除influxdb外,仓库中已按同一模式实现的插件还包括 plugins/inputs/prometheus/prometheus.go 与 plugins/inputs/influxdb_v2_listener/influxdb_v2_listener.go,可作为多类型插件(输入与服务输入)的对照参考。
四、统计在internal插件中的呈现
启用internal输入插件后,插件内部统计与模型层统计一起进入常规采集管线。该插件的完整配置如下(plugins/inputs/internal/README.md):
# Collect statistics about itself [[inputs.internal]] ## If true, collect telegraf memory stats. # collect_memstats = true ## If true, collect metrics from Go's runtime.metrics. For a full list see: ## https://pkg.go.dev/runtime/metrics # collect_gostats = false ## Collect statistics per plugin instance and not per plugin type # per_instance = false采集逻辑见 plugins/inputs/internal/internal.go:
PerInstance为false(默认)时调用collectAccumulatedPluginStat:以"测量名 + tags(剔除_id)"为 key 将同类型插件的统计聚合累加为一条指标,模拟按插件类型统计的旧行为;PerInstance为true时调用collectIndividualPluginStat:直接逐条输出每个实例的统计,保留_id标签;- 无论哪种模式,
internal_agent测量都会被补上go_version标签,所有测量统一附加version=<telegraf 版本>(internal.go)。
结合内部统计后,典型输出(节选自 plugins/inputs/internal/README.md)形如:
internal_write,output=file,host=tyrion,version=1.99.0 buffer_limit=10000i,buffer_size=0i,errors=0i,metrics_added=18i,metrics_dropped=0i,metrics_filtered=0i,metrics_rejected=0i,metrics_written=18i,startup_errors=1i,write_errors=0i,write_time_ns=636609i 1480682800000000000 internal_gather,input=internal,host=tyrion,version=1.99.0 errors=2i,gather_errors=1i,gather_time_ns=442114i,gather_timeouts=0i,metrics_gathered=19i,startup_errors=0i 1480682800000000000 internal_http_listener,address=:8186,host=tyrion,version=1.99.0 queries_received=0i,writes_received=0i,requests_received=0i,buffers_created=0i,requests_served=0i,pings_received=0i,bytes_received=0i,not_founds_served=0i,pings_served=0i,queries_served=0i,writes_served=0i 1480682800000000000 internal_mqtt_consumer,host=tyrion,version=1.99.0 messages_received=622i,payload_size=37942i 1657282270000000000其中internal_http_listener、internal_mqtt_consumer即为各自插件通过 collector 注册的内部统计测量,字段(如requests_served、messages_received、payload_size)正是那些"只有插件实例才知道"的数据点。
五、模型层统计与内部统计的边界
为帮助读者区分两类统计的职责,这里汇总模型层各 Running 类已注册的通用统计(全部通过selfstat完成,tags 由模型层构造):
| 模型类 | 测量名 | 统计字段(节选) | 源码位置 |
|---|---|---|---|
| RunningInput | internal_gather | metrics_gathered、gather_time_ns、gather_timeouts、gather_errors、startup_errors、errors | models/running_input.go |
| RunningOutput | internal_write | metrics_written、write_time_ns、write_errors、startup_errors、metrics_filtered、errors | models/running_output.go |
| RunningAggregator | internal_aggregate | metrics_pushed、metrics_filtered、metrics_dropped、push_time_ns | models/running_aggregator.go |
| Buffer | internal_write | metrics_added、metrics_written、metrics_rejected、metrics_dropped、buffer_size、buffer_limit | models/buffer.go |
| RunningParser | internal_parser | metrics_parsed、parse_time_ns | models/running_parsers.go |
| RunningSerializer | internal_serializer | metrics_serialized、bytes_serialized、serialization_time_ns | models/running_serializer.go |
这些统计由模型层在构造 Running 对象时注册,任何插件都会自动获得;而 tsd-011 定义的内部统计则是插件按需自主注册、补充模型层覆盖不到的指标。模型层与插件内部统计共用同一个selfstat注册表,因此internal插件无需区分来源即可统一采集。
此外,internal插件对外暴露的字段名称与internal测量下各字段的完整语义,可对照 plugins/inputs/internal/README.md 中internal_agent、internal_gather、internal_write的逐字段说明查阅。
六、相关规范与测试验证
tsd-011 属于 Telegraf 的 TSD(Telegraf Specification Document)体系,规范文件统一以tsd-前缀 + 递增编号命名(如tsd-001-deprecation、tsd-010-labels-and-selectors),写作要求至少包含 Objective 与 Overview 两部分,参见 docs/specs/README.md 与 docs/specs/template.md。tsd-011 即通过 Objective、Overview、Related Issues 三部分完整描述了内部统计框架的动机与设计。
实现层面可进一步通过仓库测试验证行为:models目录下的running_input_test.go(如TestRunningInputStatisticsErrorsCount)、running_output_test.go(如TestRunningOutputStatisticsErrorsCount、TestRunningOutputStatisticsWriteErrorsCount)均直接调用selfstat.Register构造相同 tags 的统计并断言计数行为,可作为理解模型层统计注册与聚合语义的参考;selfstat包的collector_test.go、selfstat_test.go则覆盖了 collector 与全局注册表的注册、去重与清理逻辑。
七、编写新插件的检查清单
基于上述规范与实现,为自研插件接入内部统计时可遵循如下清单:
- 导入依赖:
import "github.com/influxdata/telegraf/selfstat"; - 声明字段:在插件结构体中添加
Statistics *selfstat.Collector \toml:"-"``; - 注册统计:在
Init中使用s.Statistics.Register("measurement", "field", nil)或RegisterTiming(...)获取selfstat.Stat句柄并保存到插件字段; - 更新统计:在采集/写入/请求等业务路径上调用
Incr(v)或Set(v);计时类调用RegisterTiming并按需Incr单次耗时; - 清理统计:若插件生命周期内需要释放统计,可调用 collector 的
Unregister/UnregisterAll; - 验证输出:启用
[[inputs.internal]]并设置per_instance = true可逐实例观察带_id标签的内部统计,per_instance = false则观察按类型聚合后的结果。
遵循该模式,插件即可在保证"tags 与 alias 不丢失"的前提下,把自身最关键的运行数据通过标准的internal采集管线暴露给监控系统,为部署调优与故障定位提供数据支撑。
八、总结
tsd-011 规范为 Telegraf 的插件内部统计建立了一套标准机制:插件声明*selfstat.Collector类型的Statistics成员,模型层在Init之前把携带 alias 与 tags 的 collector 注入插件,插件经由 collector 注册、更新、注销自身指标,最终由internal输入插件以internal_<plugin_name>测量统一对外输出。整条链路以selfstat全局注册表为底座,以 models/common.go 的SetStatisticsOnPlugin为注入枢纽,以influxdb等真实插件为范例——理解这一设计,就能为任何 Telegraf 插件低成本地补充可观测性,也为阅读和评审其他插件实现提供了清晰的切入点。
【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考