Telegraf Burrow 输入插件:基于 HTTP API 采集 Kafka 消费滞后(Consumer Lag)的完整指南
Telegraf Burrow 输入插件基于 HTTP API 采集 Kafka 消费滞后Consumer Lag的完整指南【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf导读Burrow 是 LinkedIn 开源的 Kafka 消费滞后检查器它不依赖 ZooKeeper、不修改客户端代码即可持续评估每个消费组在每个分区上的消费进度与滞后情况。Telegraf 的inputs.burrow输入插件通过 Burrow v1.x 的 HTTP API 拉取 Kafka 集群、Topic 和消费组的实时状态将「分组状态」「分区滞后」「Topic 水位」三类数据转化为可监控、可告警的指标。读完本文你将掌握该插件的完整配置项与默认值、底层 HTTP 调用链、三类指标burrow_group/burrow_partition/burrow_topic的字段语义、状态码映射规则以及如何结合源码与测试验证采集行为。该插件自 Telegraf v1.7.0 起可用属于 messaging 类别适用于所有平台。其核心实现位于 burrow.go配置样例位于 sample.conf插件测试与 HTTP mock 数据分别位于 burrow_test.go 与 testdata。插件功能概览inputs.burrow从 Burrow 的 HTTP API默认前缀/v3/kafka读取三类信息Kafka 集群列表发现 Burrow 监控的所有集群Topic 及分区水位每个分区当前写入的 offset消费组及分区消费状态每个消费组的总滞后total lag、最大滞后maxlag以及每个分区的滞后明细、状态与 owner。插件在同一采集周期内并行处理多个服务器每个 server 一个 goroutine并在单台服务器内部对 Topic 列表与消费组列表的拉取做并发调度通过信号量guard限制并发连接数避免在 Topic 或消费组数量很大时压垮 Burrow。配置详解完整配置示例见 sample.conf所有参数均已列入下方。其中大部分参数在源码 burrow.go 中有明确的默认常量定义const ( defaultBurrowPrefix /v3/kafka defaultConcurrentConnections 20 defaultResponseTimeout time.Second * 5 defaultServer http://localhost:8000 )基础配置[[inputs.burrow]] ## Burrow API endpoints in format schema://host:port. ## Default is http://localhost:8000. servers [http://localhost:8000]serversBurrow API 端点列表格式为schema://host:port。默认值为http://localhost:8000如果该配置为空数组插件会在Gather时自动填入默认值见 burrow.go。支持配置多个端点插件会为每个 server 启动一个 goroutine 并发采集测试用例TestMultipleServers验证了多服务器场景下指标可正常聚合。api_prefix覆盖 Burrow API 前缀默认/v3/kafka。当 Burrow 部署在反向代理之后、实际路径不是默认前缀时使用。注意插件在解析 server 地址时如果 URL 本身带有路径则优先使用 URL 自带路径否则才拼接api_prefix见 burrow.go。response_timeout接收响应的最大等待时间默认5s。它同时作用于 HTTP client 的整体超时、TCP 拨号超时net.Dialer{Timeout}以及响应头超时ResponseHeaderTimeout详见createClientburrow.go。若配置值小于 1 秒会被重置为默认的 5 秒见setDefaults。concurrent_connections每台服务器允许的最大并发连接数默认20。该值通过两个渠道生效HTTP transport 的MaxIdleConnsPerHost ConcurrentConnections / 2与MaxConnsPerHost ConcurrentConnections若该值 1则空闲连接数使用 Go 默认值 2若为 0 则视为「不限制」采集过程中每个 Topic/消费组的明细请求通过带缓冲的 channelguard获取令牌从而限制并发请求数见 burrow.go。过滤器配置## Filter clusters, default is no filtering. ## Values can be specified as glob patterns. # clusters_include [] # clusters_exclude [] ## Filter consumer groups, default is no filtering. ## Values can be specified as glob patterns. # groups_include [] # groups_exclude [] ## Filter topics, default is no filtering. ## Values can be specified as glob patterns. # topics_include [] # topics_exclude []插件针对集群cluster、消费组group、主题topic各提供include/exclude两组过滤规则默认均不过滤。规则支持glob 通配模式插件通过 Telegraf 公共过滤器filter.NewIncludeExcludeFilter在首次Gather时编译这些模式见compileGlobsburrow.go。需要特别注意的是topics_include/topics_exclude的过滤同时作用于两个位置burrow_topic指标的生成以及burrow_partition分区指标因为消费组滞后明细里也带有 topic 信息见 burrow.go。认证与 TLS 配置## Credentials for basic HTTP authentication. # username # password ## Optional TLS config # tls_ca /etc/telegraf/ca.pem # tls_cert /etc/telegraf/cert.pem # tls_key /etc/telegraf/key.pem # insecure_skip_verify falseusername/passwordHTTP Basic 认证凭据。当username非空时插件在构建请求时调用req.SetBasicAuth(b.Username, b.Password)见 burrow.go。测试用例TestBasicAuthConfig通过带 Basic Auth 校验的 mock 服务器验证了该流程。tls_ca/tls_cert/tls_key/insecure_skip_verify启用 HTTPS 端点时使用的 TLS 客户端配置由插件内嵌的tls.ClientConfig提供通过b.ClientConfig.TLSConfig()生成*tls.Config后注入 HTTP transport。相关配置通用规则参见 docs/includes/plugin_config.md。除上述插件专属参数外inputs.burrow还支持 Telegraf 全局/通用插件配置如修改指标名、标签、字段、设置别名、调整插件执行顺序等详见 CONFIGURATION.md。底层采集流程与 HTTP 调用链插件每次Gather会执行如下调用链对应 burrow.go 的Gather→gatherServer→gatherTopics/gatherGroups获取集群列表GET serverapi_prefix默认即GET http://localhost:8000/v3/kafka响应体为{clusters: [clustername1, ...]}对每个集群并发拉取 Topic 列表与消费组列表GET api_prefix/cluster/topicGET api_prefix/cluster/consumer对每个 Topic 拉取分区水位GET api_prefix/cluster/topic/topic响应体为{offsets: [459178195, 459178022, ...]}数组下标即分区号据此生成burrow_topic指标对每个消费组拉取滞后状态GET api_prefix/cluster/consumer/group/lag响应体包含status组状态、partition_count、maxlag、totallag以及partitions数组每项含 topic、partition、owner、start/end offset、current_lag 等据此生成burrow_group与burrow_partition指标。在构造子资源 URL 时插件会对路径片段做url.PathEscape转义后再拼接见appendPathToURL保证集群名、组名、Topic 名中的特殊字符不会破坏 URL 结构。请求只接受 HTTP 200 响应其他状态码如 404会返回wrong response: code错误并上报到 accumulator见 burrow.go。测试中的 mock 服务器getResponseJSON将 URI 直接映射到 testdata 下的 JSON 文件例如/v3/kafka/clustername1/consumer/group1/lag对应v3_kafka_clustername1_consumer_group1_lag.json可以据此直观理解每个端点的真实响应格式。指标详解插件输出三类指标字段与标签说明如下与 README 及源码genTopicMetrics、genGroupStatusMetrics、genGroupLagMetrics一一对应。burrow_group每个消费组一条事件字段字段类型含义statusstring组状态取值见「状态映射」status_codeint状态数值码1..6见「状态映射」partition_countint分区数若响应中为 0则回退为 partitions 数组长度offsetint64所有分区end.offset之和total_lagint64总滞后Burrow 计算的totallaglagint64最大滞后maxlag.current_lag若maxlag为空则为 0timestampint64所有分区end.timestamp中的最大值标签clusterstring、groupstringburrow_partition每个主题分区一条事件字段字段类型含义statusstring分区状态见「状态映射」status_codeint状态数值码1..6见「状态映射」lagint64当前滞后current_lag为空则 0offsetint64end.offsettimestampint64end.timestamp标签cluster、group、topicstring、partitionint、ownerstringburrow_topic每个 Topic 分区一条事件字段offsetint64该分区当前写入水位标签cluster、topicstring、partitionint由 offsets 数组下标转换而来注意 README 中对burrow_partition.offset的描述写作end.timestamp但从源码 burrow.go 与测试 burrow_test.go 的实际断言看该字段取的是partition.End.Offsetint64 类型读者以源码实现为准。状态映射规则Burrow 返回的状态字符串会被转换为数值码便于在时序数据库中进行数值比较与告警状态数值码OK1NOT_FOUND2WARN3ERR4STOP5STALL6未知值0该映射在mapStatusToCode中实现burrow.gostatus_code字段直接引用该函数结果。0意味着「未知状态」在告警规则中应予以区分处理。测试与验证插件测试 burrow_test.go 使用httptest构造 mock 服务器按 URI 返回 testdata 中对应的 JSON 样本覆盖了以下关键场景TestBurrowTopic/TestBurrowPartition/TestBurrowGroup验证三类指标的字段与标签映射TestMultipleServers验证多端点并发采集TestMultipleRuns验证重复采集的稳定性4 次采集每次均产出 7 条指标TestBasicAuthConfig验证 Basic Auth 认证流程TestFilterClusters/TestFilterGroups/TestFilterTopics验证 glob 过滤器的 include/exclude 语义例如group?匹配group1*全排除。参考数据文件可以让你快速理解 Burrow 的真实响应结构集群列表见v3_kafka.jsonTopic 水位见v3_kafka_clustername1_topic_topicA.jsonoffsets 数组消费组滞后状态见v3_kafka_clustername1_consumer_group1_lag.json包含status、partition_count、maxlag、totallag与partitions明细。典型使用场景与示例输出基于 testdata 中的样本clustername1集群、topicA主题、group1消费组、3 个分区、状态全部OK一次采集会生成类似如下的指标时间为 Burrow 返回的毫秒时间戳 1515609490008即 2018-01-10burrow_group,clusterclustername1,groupgroup1 statusOK,status_code1i,partition_count3i,offset1291282720i,total_lag0i,lag0i,timestamp1515609490008i burrow_partition,clusterclustername1,groupgroup1,topictopicA,partition0,ownerkafka1 statusOK,status_code1i,lag0i,offset431323195i,timestamp1515609490008i burrow_partition,clusterclustername1,groupgroup1,topictopicA,partition1,ownerkafka2 statusOK,status_code1i,lag0i,offset431322962i,timestamp1515609490008i burrow_partition,clusterclustername1,groupgroup1,topictopicA,partition2,ownerkafka3 statusOK,status_code1i,lag0i,offset428636563i,timestamp1515609490008i burrow_topic,clusterclustername1,topictopicA,partition0 offset459178195i burrow_topic,clusterclustername1,topictopicA,partition1 offset459178022i burrow_topic,clusterclustername1,topictopicA,partition2 offset456491598i实际落地后常见的告警思路是对burrow_group的status_code 3即出现 WARN/ERR/STOP/STALL进行告警对lag/total_lag设置阈值在消费滞后持续增长时触发告警对burrow_topic.offset可监控写入水位停滞辅助判断生产端是否异常。结合clusters_include、groups_include与topics_include过滤可以只监控关心的集群、消费组和主题降低采集压力与指标基数。【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →