ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Grafana Tempo 中 kprom 插件详解:为 franz-go Kafka 客户端注入 Prometheus 指标

Grafana Tempo 中 kprom 插件详解:为 franz-go Kafka 客户端注入 Prometheus 指标 Grafana Tempo 中 kprom 插件详解为 franz-go Kafka 客户端注入 Prometheus 指标【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempokprom 是 franz-gokgoKafka 客户端库提供的 Prometheus 指标插件包位于本仓库的 vendor 目录 kprom 包 中。它通过实现kgo.Hook接口在不侵入客户端业务逻辑的前提下把 Kafka 客户端的连接、读写、生产/消费等关键数据自动暴露为 Prometheus 指标。在 Tempo 中distributor 写 Kafka、ingest storage reader 读 Kafka、livestore 分区消费者等 Kafka 交互链路都依赖该插件来构建可观测性。读完本文你将掌握 kprom 的指标清单、全部配置选项registry、直方图桶、fetch/produce 细节标签并了解 Tempo 在 pkg/ingest/reader_client.go 与 pkg/ingest/writer_client.go 中如何实际组装和使用它。基础用法一行 Hook 挂上整套指标kprom 的使用方式极简创建Metrics后通过kgo.WithHooks注入客户端即可。这正是 README 中给出的最小示例metrics : kprom.NewMetrics(namespace) cl, err : kgo.NewClient( kgo.WithHooks(metrics), // ...other opts )NewMetrics(namespace, opts...)接收一个指标命名空间namespace和若干Opt配置项返回的*Metrics同时实现了 kgo 的多组 Hook 接口。从 kprom.go 顶部的接口断言可以看到它挂载的完整 Hook 面var ( // interface checks to ensure we implement the hooks properly _ kgo.HookBrokerConnect new(Metrics) _ kgo.HookBrokerDisconnect new(Metrics) _ kgo.HookBrokerWrite new(Metrics) _ kgo.HookBrokerRead new(Metrics) _ kgo.HookProduceBatchWritten new(Metrics) _ kgo.HookFetchBatchRead new(Metrics) _ kgo.HookBrokerE2E new(Metrics) _ kgo.HookBrokerThrottle new(Metrics) _ kgo.HookNewClient new(Metrics) _ kgo.HookClientClosed new(Metrics) )两个关键的生命周期行为值得注意见 OnNewClient / OnClientClosed指标是“懒注册”的所有 Counter/Histogram 都在OnNewClient回调中真正创建并注册到 registry而不是在NewMetrics时。这意味着同一个*Metrics可以为多个客户端依次注册并且当某个客户端Close()时OnClientClosed会把该客户端注册的所有 collector 从 registry 中Unregister避免指标残留。支持作为外部自定义 CollectorMetrics本身实现了Collect/Describekprom.go#L512-L528配合kprom.Registry(nil)可以把 kprom 整体作为一个自定义 prometheus collector 挂到外部 registry 上实现多客户端共用一个 registry。此外包还提供两个便捷访问器Metrics.Registry()返回内部 registry便于你继续往里加自己的指标Metrics.Handler()直接返回promhttp的http.Handler可以挂到自己的 HTTP 服务上供 Prometheus 抓取。默认指标清单10 个计数器与默认标签README 列出了默认无任何扩展选项时暴露的一组计数器counter vec指标全部带命名空间前缀#{ns}#{ns}_connects_total{node_id#{node}} #{ns}_connect_errors_total{node_id#{node}} #{ns}_write_errors_total{node_id#{node}} #{ns}_write_bytes_total{node_id#{node}} #{ns}_read_errors_total{node_id#{node}} #{ns}_read_bytes_total{node_id#{node}} #{ns}_produce_bytes_total{node_id#{node},topic#{topic}} #{ns}_fetch_bytes_total{node_id#{node},topic#{topic}} #{ns}_buffered_produce_records_total #{ns}_buffered_fetch_records_total各指标含义与触发点结合 kprom.go 的 Hook 实现指标名类型标签说明#{ns}_connects_totalCounternode_id成功建立的连接总数在OnBrokerConnect且err nil时自增#{ns}_connect_errors_totalCounternode_id连接失败次数OnBrokerConnect收到err时自增#{ns}_write_errors_totalCounternode_id写请求错误次数由OnBrokerE2E中e2e.WriteErr ! nil触发#{ns}_write_bytes_totalCounternode_id写入 TCP 连接的字节数压缩后的字节OnBrokerE2E中累加e2e.BytesWritten#{ns}_read_errors_totalCounternode_id读请求错误次数OnBrokerE2E中e2e.ReadErr ! nil时触发#{ns}_read_bytes_totalCounternode_id从 TCP 连接读取的字节数解压前的字节累加e2e.BytesRead#{ns}_produce_bytes_totalCounternode_id, topic生产的未压缩字节数OnProduceBatchWritten触发#{ns}_fetch_bytes_totalCounternode_id, topic消费的未压缩字节数OnFetchBatchRead触发#{ns}_buffered_produce_records_totalGauge—客户端内待发送的缓冲记录数通过GaugeFunc直接读client.BufferedProduceRecords()#{ns}_buffered_fetch_records_totalGauge—客户端内待消费的缓冲记录数读client.BufferedFetchRecords()两点补充说明node_id 标签的 seed 前缀约定README 特别指出seed broker种子节点的 ID 会带seed_前缀数字表示是第几个 seed如seed_0、seed_1。这意味着在监控面板里按node_id分组时seed_N并不是真正的 broker 节点排查连接问题时需要注意区分。缓冲指标是“采样式”Gauge从源码结构看两个 buffered 指标并非每次 Hook 回调时更新而是注册为prometheus.GaugeFunckprom.go#L347-L367在 Prometheus 每次抓取时才实时查询客户端当前缓冲量因此它反映的是抓取瞬间的快照值。扩展指标直方图计时与 fetch/produce 细节README 提到“上述指标可以通过本包选项大幅扩展包括计时、未压缩/压缩字节、不同标签”对应 config.go 中的两类选项。1. 启用直方图计时指标默认情况下所有直方图都不启用。可通过Histograms(...)开启若干个计时直方图使用默认桶或用HistogramsFromOpts(...)为每个直方图单独指定桶metrics : kprom.NewMetrics( kprom.Histograms(kprom.RequestDurationE2E), )可启用的Histogram常量共 6 个config.go#L141-L151以及它们记录的语义常量指标名语义ReadWait#{ns}_read_wait_seconds等待从 Kafka 读响应的时间ReadTime#{ns}_read_time_seconds实际读取耗时WriteWait#{ns}_write_wait_seconds等待写入 Kafka 的时间WriteTime#{ns}_write_time_seconds实际写入耗时RequestDurationE2E#{ns}_request_duration_e2e_seconds请求写入开始到响应完全读完的端到端时间RequestThrottled#{ns}_request_throttled_seconds被 broker 节流throttle的时长所有直方图默认使用DefBucketsvar DefBuckets []float64{0.001, 0.002, 0.004, 0.008, 0.016, 0.032, 0.064, 0.128, 0.256, 0.512, 1.024, 2.048}这是一组按 2 倍递增、覆盖 1ms2s 的桶注释说明是“为 Kafka 常见时延量级量身定制”。可以用Buckets([]float64)全局覆盖或用HistogramsFromOpts精细控制metrics : kprom.NewMetrics( kprom.HistogramsFromOpts( kprom.HistogramOpts{ Enable: kprom.ReadWait, Buckets: prometheus.LinearBuckets(10, 10, 8), }, kprom.HistogramOpts{ Enable: kprom.ReadTime, // 不指定 Buckets 时使用 kprom 默认桶 }, ), )一个容易忽略的实现细节直方图只有在“启用”时才会真正Observe。例如 OnBrokerE2E 中每次都先if _, ok : m.cfg.histograms[WriteWait]; ok判断后才记录未启用的直方图不产生任何观测开销。2. 配置 fetch/produce 指标的细节DetailFetchAndProduceDetail(details ...Detail)控制 produce/fetch 类指标的标签与口径默认值为UncompressedBytes ByTopic ByNode即node_id、topic标签 未压缩字节计数对应上面默认清单中的produce_bytes_total/fetch_bytes_total。可选的Detail常量config.go#L199-L210常量作用ByNodefetch/produce 指标加node_id标签ByTopicfetch/produce 指标加topic标签Batches额外报告#{ns}_produce_batches_total/#{ns}_fetch_batches_totalRecords额外报告#{ns}_produce_records_total/#{ns}_fetch_records_totalCompressedBytes额外报告压缩字节数#{ns}_produce_compressed_bytes_total/#{ns}_fetch_compressed_bytes_totalUncompressedBytes报告未压缩字节数默认指标名见下ConsistentNaming将produce_bytes_total/fetch_bytes_total重命名为produce_uncompressed_bytes_total/fetch_uncompressed_bytes_total与压缩字节指标命名对齐注意FetchAndProduceDetail是全量覆盖而非增量追加内部实现会先把fetchProduceOpts重置为空结构再逐项设置config.go#L223-L253只传Batches就会丢掉默认的字节计数。Registry、标签与其他选项config.go 还定义了以下配置项默认行为是“创建全新的独立 prometheus registry”newCfg中prometheus.NewRegistry()常用选项如下选项说明Registry(rg RegistererGatherer)同时指定 registerer 与 gatherer传nil表示把 kprom 作为自定义 collector 挂到外部 registry此时与GoCollectors互斥Registerer(reg)仅指定注册目标Tempo 实际最常用的方式Gatherer(g)仅指定抓取源GoCollectors()额外注册GoCollector和ProcessCollector让/metrics端点自带运行时指标WithClientLabel()给所有指标追加client_id常量标签取客户端的kgo.ClientIDWithStaticLabel(labels)给所有指标追加静态常量标签内部会maps.Clone深拷贝Subsystem(ss)设置 subsystem指标名变为#{ns}_#{ss}_...Buckets(...)/Histograms(...)/HistogramsFromOpts(...)直方图桶与开关见上一节HandlerOpts(opts)自定义Metrics.Handler()返回的 promhttp 处理器选项适合不用自带 registry 又想覆盖默认处理行为时Tempo 中的真实用法Tempo 仓库在多处直接消费 kprom可以对照源码验证上述用法。Distributor 写 Kafkapkg/ingest/writer_client.go 中NewWriterClient创建无命名空间的Metrics并把注册目标交给调用方包装好的 registerer同时全量启用四类 fetch/produce 细节func NewWriterClient(kafkaCfg KafkaConfig, maxInflightProduceRequests int, logger log.Logger, reg prometheus.Registerer) (*kgo.Client, error) { // Do not export the client ID, because we use it to specify options to the backend. metrics : kprom.NewMetrics( , // No prefix. We expect the input prometheus.Registered to be wrapped with a prefix. kprom.Registerer(reg), kprom.FetchAndProduceDetail(kprom.Batches, kprom.Records, kprom.CompressedBytes, kprom.UncompressedBytes), ) ... }调用方 modules/distributor/distributor.go 传入prometheus.WrapRegistererWithPrefix(tempo_distributor_, reg)最终指标呈现为tempo_distributor_produce_records_total{node_id...,topic...}这样的形态。代码注释明确解释了“为什么不导出 client ID”——client_id被用来向后端传递选项导出会造成信息泄露式耦合。Ingest Storage Reader 读 Kafkapkg/ingest/reader_client.go 提供了统一工厂函数命名空间固定为tempo_ingest_storage_reader并用WrapRegistererWith追加component标签func NewReaderClientMetrics(component string, reg prometheus.Registerer) *kprom.Metrics { return kprom.NewMetrics(tempo_ingest_storage_reader, kprom.Registerer(prometheus.WrapRegistererWith(prometheus.Labels{component: component}, reg)), // Do not export the client ID, because we use it to specify options to the backend. kprom.FetchAndProduceDetail(kprom.Batches, kprom.Records, kprom.CompressedBytes, kprom.UncompressedBytes)) }Livestore 分区消费者modules/livestore/partition_reader.go 为每个分区单独创建一组kprom.Metrics并通过静态包装追加partition标签从而在 Prometheus 中按分区维度区分消费速率。多客户端共用的测试做法pkg/ingest/writer_client_test.go 与 pkg/ingest/partition_offset_client_test.go 中用kprom.NewMetrics(, kprom.Registerer(prometheus.NewPedanticRegistry()))创建严格的 pedantic registry 来验证重复注册问题这也是“客户端关闭时OnClientClosed反注册”这一机制存在的原因之一。从源码结构看Tempo 并未启用 kprom 的直方图计时选项各处NewMetrics均未传Histograms说明当前 Tempo 对 Kafka 链路的监控重点是字节/批次/记录量与缓冲水位而非请求级时延分布如需时延观测只需在相应NewMetrics调用中追加kprom.Histograms(kprom.RequestDurationE2E)即可开启。小结kprom 的价值在于用一组 Hook 就把kgo客户端变成了“自带仪表盘”默认 10 个计数器覆盖连接、TCP 读写字节、生产/消费字节与缓冲水位再通过Histograms/FetchAndProduceDetail等选项按需扩展计时与压缩字节细节通过Registerer/WithStaticLabel/WrapRegistererWithPrefix的组合融入既有 registry 的命名与标签体系。Tempo 在 distributor、ingest storage reader、livestore 分区消费者三处 Kafka 入口的用法见上文源码路径是阅读该插件最直接的真实参照。【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempo创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进