Telegraf cloud_pubsub_push 输入插件实战搭建 Google Cloud Pub/Sub Push 推送端点【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf本篇指南围绕 Telegraf 官方cloud_pubsub_push输入插件展开讲解如何让 Telegraf 作为一个 HTTP 服务端点接收 Google Cloud Pub/Sub Push 订阅推送过来的 JSON 消息并完成鉴权、解析、指标采集与可靠确认的完整链路。读完本文你将掌握该插件的全部配置参数、其基于 tracking metrics 的可靠投递原理以及如何结合仓库源码与测试用例验证其行为。插件是什么把 Telegraf 变成 Pub/Sub 的 Push 端点cloud_pubsub_push是 Telegraf 的一个服务型输入插件Service Input自Telegraf v1.10.0起提供参见 插件 README。与常规输入插件按固定interval主动拉取数据不同它启动一个 HTTP 监听服务被动等待 Google Cloud Pub/Sub 通过 HTTP POST 推送消息。该插件接收的是 Google Pub/Sub 的JSON 格式推送消息Pub/Sub Push Subscription 的消息体其中消息data字段以 Base64 编码承载实际负载。典型应用场景包括物联网设备遥测数据经由 Cloud IoT Core → Pub/Sub → Push 订阅直接落入 Telegraf将 Pub/Sub 消息流接入 Telegraf 的解析、聚合与输出管道统一汇入 InfluxDB 等存储在无法使用 Pub/Sub Pull 订阅或需要外部服务主动回调时作为轻量级接收端点。从插件注册代码可以看到其声明方式cloud_pubsub_push.goinputs.Add(cloud_pubsub_push, func() telegraf.Input { return PubSubPush{ ServiceAddress: :8080, Path: /, MaxUndeliveredMessages: defaultMaxUndeliveredMessages, } })即插件默认监听:8080、路径为/、最大未投递消息数为1000。前置条件与安全约束使用该插件前需要理解 Google Pub/Sub Push 的一个关键行为Pub/Sub 服务只通过 HTTPS/TLS 发送推送。因此插件必须满足以下任一条件参见 插件 README部署在一个可终止 TLS 的反向代理如 Nginx、负载均衡器之后将明文 HTTP 转发给 Telegraf或直接为插件配置 TLS通过tls_cert与tls_key提供服务端证书与私钥如需双向认证mTLS可进一步在tls_allowed_cacerts中列出允许的客户端 CA 证书文件从而只授权携带受信证书的客户端连接。另外插件本身提供了简单的应用级鉴权设置token后会要求请求携带与之一致的 token 参数详见下文请求处理链路。推送订阅Push Subscription需要预先在 Pub/Sub 中为某个 Topic 创建插件只负责接收不会替你创建订阅。服务型输入的特点该插件属于服务型输入Service Input其文档明确了两点与普通插件的差异参见 service_input.md全局或插件级interval设置对它可能不生效——它依赖外部推送驱动而不是定时采集CLI 选项--test、--test-wait、--once可能不会为该插件产生输出——因为这些模式面向一次性/定时采集而服务型插件需要等待外部事件。这意味着调试时应采用先启动、再手动 POST 测试消息的方式而不是依赖--test查看输出。配置详解一个完整的示例配置插件的完整示例配置见 sample.conf也是 README 中嵌入的sample.conf。以下是带注释的完整配置# Google Cloud Pub/Sub Push HTTP listener [[inputs.cloud_pubsub_push]] ## Address and port to host HTTP listener on service_address :8080 ## Application secret to verify messages originate from Cloud Pub/Sub # token ## Path to listen to. # path / ## Maximum duration before timing out read of the request # read_timeout 10s ## Maximum duration before timing out write of the response. This should be ## set to a value large enough that you can send at least metric_batch_size ## number of messages within the duration. # write_timeout 10s ## Maximum allowed http request body size in bytes. ## 0 means to use the default of 524,288,00 bytes (500 mebibytes) # max_body_size 500MB ## Whether to add the pubsub metadata, such as message attributes and ## subscription as a tag. # add_meta false ## Max undelivered messages ## This plugin uses tracking metrics, which ensure messages are read to ## outputs before acknowledging them to the original broker to ensure data ## is not lost. This option sets the maximum messages to read from the ## broker that have not been written by an output. ## ## This value needs to be picked with awareness of the agents ## metric_batch_size value as well. Setting max undelivered messages too high ## can result in a constant stream of data batches to the output. While ## setting it too low may never flush the brokers messages. # max_undelivered_messages 1000 ## Set one or more allowed client CA certificate file names to ## enable mutually authenticated TLS connections # tls_allowed_cacerts [/etc/telegraf/clientca.pem] ## Add service certificate and key # tls_cert /etc/telegraf/cert.pem # tls_key /etc/telegraf/key.pem ## Data format to consume. ## Each data format has its own unique set of configuration options, read ## more about them here: ## https://github.com/influxdata/telegraf/blob/master/docs/DATA_FORMATS_INPUT.md data_format influx各参数说明如下参数默认值说明service_address:8080HTTP 监听地址与端口:8080表示监听所有网卡的 8080 端口token空不校验应用层共享密钥用于验证消息确实来自 Cloud Pub/Subpath/监听路径Pub/Sub 推送会打到service_address pathread_timeout10s请求读取超时若配置值小于 1 秒源码会强制回退为 10 秒见下文write_timeout10s响应写入超时应足够长以便在一次超时窗口内至少投递metric_batch_size条消息max_body_size500MB500 * 1024 * 1024字节允许的最大请求体大小超限返回 HTTP 413设为0则使用默认值add_metafalse是否把 Pub/Sub 消息属性attributes与订阅名subscription作为 tag 写入指标max_undelivered_messages1000已读但尚未被输出写入的最大消息数配合 tracking metricstls_allowed_cacerts空允许的客户端 CA 证书列表用于启用双向 TLStls_cert/tls_key空服务端证书与私钥用于启用 TLSdata_formatinflux消息data字段解析所用的数据格式关于max_body_size与超时参数的默认回退逻辑源码中体现为cloud_pubsub_push.goconst ( // 500 MB defaultMaxBodySize 500 * 1024 * 1024 defaultMaxUndeliveredMessages 1000 )if p.MaxBodySize 0 { p.MaxBodySize config.Size(defaultMaxBodySize) } if p.ReadTimeout config.Duration(time.Second) { p.ReadTimeout config.Duration(time.Second * 10) } if p.WriteTimeout config.Duration(time.Second) { p.WriteTimeout config.Duration(time.Second * 10) }也就是说即使你在配置里写了read_timeout 100ms低于 1 秒的值也会被强制提升为 10 秒避免过短的超时导致请求频繁失败。data_format消息负载如何被解析data_format决定 Pub/Sub 消息中data字段Base64 解码后被解析成 Telegraf 指标的方式默认是 Influx 行协议。该字段引用的完整数据格式清单见 DATA_FORMATS_INPUT.md。常见的选项包括influx默认负载为 InfluxDB 行协议文本json负载为 JSON可配合tag_keys等子选项提取字段与标签value负载为单值需配合data_typecsv、graphite、nagios等其余 Telegraf 支持的输入格式。解析通过SetParser注入到插件中telegraf.Parser接口插件本身不关心具体格式只负责把 Base64 数据解码后交给解析器。请求处理链路从 HTTP 请求到指标确认深入源码可以还原一条完整的处理链路。插件将PubSubPush自身注册为http.Server的 Handlercloud_pubsub_push.go核心流程如下1. 启动服务Startp.server http.Server{ Addr: p.ServiceAddress, Handler: http.TimeoutHandler(p, time.Duration(p.WriteTimeout), timed out processing metric), ReadTimeout: time.Duration(p.ReadTimeout), TLSConfig: tlsConf, }Start中会通过p.ServerConfig.TLSConfig()构造 TLS 配置若tlsConf非空则调用ListenAndServeTLS(, )证书直接由配置注入否则调用ListenAndServe()。同时它创建了 tracking accumulatorp.acc acc.WithTracking(p.MaxUndeliveredMessages) p.sem make(chan struct{}, p.MaxUndeliveredMessages) p.undelivered make(map[telegraf.TrackingID]chan bool)并启动一个后台 goroutinereceiveDelivered()持续监听已投递通知。2. 路由与鉴权ServeHTTP/authenticateIfSet请求路径必须与配置的path完全一致否则返回 404若设置了token会通过常量时间比较subtle.ConstantTimeCompare校验req.FormValue(token)失败返回 HTTP 401——用常量时间比较可避免时序侧信道攻击。3. 校验与方法检查serveWrite若请求上下文已取消或缓冲区已满信号量p.sem取不到令牌返回 HTTP 503ContentLength超过max_body_size返回 HTTP 413StatusRequestEntityTooLarge非 POST 方法返回 HTTP 405StatusMethodNotAllowed使用http.MaxBytesReader限制读取体大小读取失败同样返回 413。4. 解码与解析var payload payload if err json.Unmarshal(bytes, payload); err ! nil { ... 400 ... } sDec, err : base64.StdEncoding.DecodeString(payload.Msg.Data) if err ! nil { ... 400 ... } metrics, err : p.Parse(sDec) if err ! nil { ... 400 ... }先将请求体解析为 Pub/Sub 推送结构体payload含message.data、message.attributes、subscription再对data做 Base64 解码最后交给p.Parse按data_format解析成指标任何一步失败均返回 HTTP 400。5. 元数据注入可选当add_meta true时把消息的全部 attributes 作为 tag 加入指标并额外加入subscription标签值为projects/.../subscriptions/...if p.AddMeta { for i : range metrics { for k, v : range payload.Msg.Atts { metrics[i].AddTag(k, v) } metrics[i].AddTag(subscription, payload.Subscription) } }6. 追踪与确认tracking metricsch : make(chan bool, 1) p.mu.Lock() p.undelivered[p.acc.AddTrackingMetricGroup(metrics)] ch p.mu.Unlock()插件把这一组指标注册为tracking metric group然后阻塞等待后台receiveDelivered()通过p.acc.Delivered()返回的投递结果指标成功送达所有输出 → 向 Pub/Sub 返回HTTP 204No Content表示可以确认ack消息投递失败如输出写入报错→ 返回HTTP 500Pub/Sub 会按重试策略重新推送等待期间若请求上下文取消 → 返回 503。为什么需要max_undelivered_messages这一机制对应 Telegraf 的tracking metrics特性参见 METRICS.md 的 Tracking Metrics 一节输入在把消息交给输出之前不会向消息源确认。若 Telegraf 或系统在投递过程中崩溃未送达的消息可以在之后被 Pub/Sub 重新推送从而保证数据不丢失。max_undelivered_messages就是已读取但尚未被输出写走的消息上限它同时控制信号量p.sem的容量。配置时需要注意与 agent 的metric_batch_size配合设置过高可能导致 Telegraf 持续不断地把数据批次推向输出冲击刷新间隔与下游设置过低信号量很快占满插件对外表现为 503Pub/Sub 消息可能迟迟无法被消费。其核心权衡逻辑可以参见 METRICS.md 对 undelivered messages 的说明。测试用例插件行为的可验证依据插件行为在 cloud_pubsub_push_test.go 中有完整的表驱动测试覆盖可以当作行为契约来读。测试用例覆盖了以下场景与对应状态码测试场景期望 HTTP 状态码GET 方法访问405 Method Not AllowedPOST 到未配置的路径404 Not Found请求体超过大小限制413 Request Entity Too Large正常合法 POSTdata为正确 Base64 的 Influx 行协议204 No Content输出写入失败failWritetrue503/500 投递失败路径信号量缓冲区已满fulltrue503 Service Unavailable请求体非法 JSON400 Bad RequestBase64 解码失败400 Bad Request解码后数据格式非法400 Bad Request测试中使用的消息体是典型的 Pub/Sub Push 格式例如含 attributes、Base64 data、subscription{message:{attributes:{deviceId:myPi,deviceNumId:2808946627307959,deviceRegistryId:my-registry,deviceRegistryLocation:us-central1,projectId:conference-demos,subFolder:},data:dGVzdGluZ0dvb2dsZSxzZW5zb3I9Ym1lXzI4MCB0ZW1wX2M9MjMuOTUsaHVtaWRpdHk9NjIuODMgMTUzNjk1Mjk3NDU1MzUxMDIzMQ,messageId:204004313210337,publishTime:2018-09-14T19:22:54.587Z},subscription:projects/conference-demos/subscriptions/my-subscription}其中data解码后正是 Influx 行协议文本testingGoogle,sensorbme_280 temp_c23.95,humidity62.83 1536952974553510231。你可以用同样的结构手动构造 POST 请求来联调自己的 Telegraf 实例curl -i -X POST \ -H Content-Type: application/json \ -d {message:{attributes:{},data:dGVzdGluZ0dvb2dsZSxzZW5zb3I9Ym1lXzI4MCB0ZW1wX2M9MjMuOTUsaHVtaWRpdHk9NjIuODMgMTUzNjk1Mjk3NDU1MzUxMDIzMQ,messageId:test-1},subscription:projects/demo/subscriptions/test} http://localhost:8080/成功时返回 204并可在 Telegraf 的输出端看到testingGoogle测量。部署与运维要点网络拓扑由于 Pub/Sub 只走 HTTPS生产环境推荐在 Telegraf 前置反向代理终止 TLSTelegraf 侧监听内网地址也可以直接在插件上配置tls_cert/tls_key直连公网。启用tls_allowed_cacerts可强制客户端证书验证形成 mTLS。调试建议不要在--test模式下期望输出先正常启动 agent再用上面的 curl 模拟 Pub/Sub 推送。配合data_format与日志观察 4xx/5xx 响应即可定位问题。容量规划关注max_undelivered_messages与 agent 的metric_batch_size、flush_interval的关系避免持续满负荷推流或消息积压max_body_size与代理侧上传大小限制需要保持一致。消息确认语义204 即代表已送达输出若输出失败返回 5xxPub/Sub 会按订阅重试策略再次推送。这是该插件数据不丢失的关键保证因此下游输出的稳定性直接影响消息消费进度。小结cloud_pubsub_push让 Telegraf 在无需任何额外组件的情况下成为 Google Cloud Pub/Sub Push 订阅的标准 HTTP 端点。它把 Pub/Sub 的 JSON 推送格式、Base64 解码、可插拔的data_format解析、可选的元数据标签、token 与 mTLS 双重鉴权以及基于 tracking metrics 的可靠确认机制整合在一起。理解max_undelivered_messages与 204/500 响应语义是把它稳定接入生产管道、做到端到端不丢数据的关键。若需要进一步了解插件全局配置与数据处理选项可继续查阅 CONFIGURATION.md 的 Plugins 一节 与 DATA_FORMATS_INPUT.md。【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考