Telegraf 插件内部统计(internal plugin statistics):基于 selfstat.Collector 的插件级指标采集规范与实践
2026/9/13 12:04:10 网站建设 项目流程

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改为逐实例输出。

其中"插件统计"又分为两类:

  1. 模型层(model-level)统计:由models包中的各 Running 类(如running_input.gorunning_output.go)负责注册,覆盖采集时长、写入时长、错误次数等通用维度;
  2. 插件内部(plugin-internal)统计:只有插件实例自身知道的信息,例如向输出目标实际写入的字节数、特定协议错误、请求计数等,无法在模型层获知。

问题在于:插件若绕过模型层直接调用selfstat.Register注册统计对象,将无法获得全局tags配置与alias别名设置——这些信息只在模型层掌握。tsd-011 规范的初衷正是解决这一缺陷:定义一个框架,把携带了alias与 tags 信息的"统计收集器"(statistics collector)注入插件,让插件通过它注册内部统计。相关背景可追溯至仓库历史中的 issue #4889(InfluxDB 输出插件内部统计)、#6965(Kafka 输出插件内部统计)与 #17275(InfluxDB v2 输出插件内部统计)。

二、总体机制:注入一个统计收集器

规范提出的核心设计包含三个环节:

  1. 插件侧:插件结构体导出Statistics成员,类型必须为指针*selfstat.Collector,强烈建议为其声明"-"的 TOML tag,避免用户配置与内部成员冲突;
  2. 模型层注入:Telegraf 模型代码在实例化插件之后、调用插件Init函数之前,将selfstat.Collector实例注入Statistics成员,注入时携带模型层已知的aliastags等全部信息;
  3. 插件使用:插件必须通过该 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) StatReset(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_idalias等模型层标签。

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:

  • PerInstancefalse(默认)时调用collectAccumulatedPluginStat:以"测量名 + tags(剔除_id)"为 key 将同类型插件的统计聚合累加为一条指标,模拟按插件类型统计的旧行为;
  • PerInstancetrue时调用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_listenerinternal_mqtt_consumer即为各自插件通过 collector 注册的内部统计测量,字段(如requests_servedmessages_receivedpayload_size)正是那些"只有插件实例才知道"的数据点。

五、模型层统计与内部统计的边界

为帮助读者区分两类统计的职责,这里汇总模型层各 Running 类已注册的通用统计(全部通过selfstat完成,tags 由模型层构造):

模型类测量名统计字段(节选)源码位置
RunningInputinternal_gathermetrics_gatheredgather_time_nsgather_timeoutsgather_errorsstartup_errorserrorsmodels/running_input.go
RunningOutputinternal_writemetrics_writtenwrite_time_nswrite_errorsstartup_errorsmetrics_filterederrorsmodels/running_output.go
RunningAggregatorinternal_aggregatemetrics_pushedmetrics_filteredmetrics_droppedpush_time_nsmodels/running_aggregator.go
Bufferinternal_writemetrics_addedmetrics_writtenmetrics_rejectedmetrics_droppedbuffer_sizebuffer_limitmodels/buffer.go
RunningParserinternal_parsermetrics_parsedparse_time_nsmodels/running_parsers.go
RunningSerializerinternal_serializermetrics_serializedbytes_serializedserialization_time_nsmodels/running_serializer.go

这些统计由模型层在构造 Running 对象时注册,任何插件都会自动获得;而 tsd-011 定义的内部统计则是插件按需自主注册、补充模型层覆盖不到的指标。模型层与插件内部统计共用同一个selfstat注册表,因此internal插件无需区分来源即可统一采集。

此外,internal插件对外暴露的字段名称与internal测量下各字段的完整语义,可对照 plugins/inputs/internal/README.md 中internal_agentinternal_gatherinternal_write的逐字段说明查阅。

六、相关规范与测试验证

tsd-011 属于 Telegraf 的 TSD(Telegraf Specification Document)体系,规范文件统一以tsd-前缀 + 递增编号命名(如tsd-001-deprecationtsd-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(如TestRunningOutputStatisticsErrorsCountTestRunningOutputStatisticsWriteErrorsCount)均直接调用selfstat.Register构造相同 tags 的统计并断言计数行为,可作为理解模型层统计注册与聚合语义的参考;selfstat包的collector_test.goselfstat_test.go则覆盖了 collector 与全局注册表的注册、去重与清理逻辑。

七、编写新插件的检查清单

基于上述规范与实现,为自研插件接入内部统计时可遵循如下清单:

  1. 导入依赖import "github.com/influxdata/telegraf/selfstat"
  2. 声明字段:在插件结构体中添加Statistics *selfstat.Collector \toml:"-"``;
  3. 注册统计:在Init中使用s.Statistics.Register("measurement", "field", nil)RegisterTiming(...)获取selfstat.Stat句柄并保存到插件字段;
  4. 更新统计:在采集/写入/请求等业务路径上调用Incr(v)Set(v);计时类调用RegisterTiming并按需Incr单次耗时;
  5. 清理统计:若插件生命周期内需要释放统计,可调用 collector 的Unregister/UnregisterAll
  6. 验证输出:启用[[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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询