news 2026/9/15 1:43:54

Telegraf CloudWatch Metric Streams 输入插件实战指南:基于 Firehose HTTP 交付的 AWS 指标流接入与配置

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Telegraf CloudWatch Metric Streams 输入插件实战指南:基于 Firehose HTTP 交付的 AWS 指标流接入与配置

Telegraf CloudWatch Metric Streams 输入插件实战指南:基于 Firehose HTTP 交付的 AWS 指标流接入与配置

【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf

Amazon CloudWatch Metric Streams 允许用户把 CloudWatch 指标以近乎实时的方式持续推送到指定的接收端(如 Kinesis Data Firehose),从而替代"按固定间隔轮询 GetMetricData API"的传统拉取模式。Telegraf 仓库中的cloudwatch_metric_streams输入插件(自 Telegraf v1.24.0 起引入)正是为此设计的服务端接收器(HTTP Listener):它监听并接收 Firehose 按 HTTP 交付规范推送过来的指标数据,解码后写入 Telegraf 指标管线,供后续处理器(processors)、聚合器(aggregators)与输出插件(outputs)使用。

阅读本篇技术指南后,你将掌握该插件的完整工作原理、全部配置项及其默认行为、收到数据后的解码与命名转换规则(含api_compatabilityAPI 兼容模式)、自监控指标与排障方法,并能基于仓库源码与测试样例搭建一套可复现的接入方案。

一、插件定位:面向 HTTP 交付的 Service Input

与大多数按interval定时抓取的输入插件不同,cloudwatch_metric_streams属于服务输入插件(Service Input)。它不主动去 AWS 拉数据,而是在本地启动一个常驻 HTTP 服务,等待 AWS Firehose 将指标批量推送过来。关于服务输入插件的通用说明可参考 docs/includes/service_input.md,其关键差异有两点:

  1. 全局或插件级interval设置可能不生效——数据的到达节奏由上游 Firehose 推送频率决定;
  2. --test--test-wait--once等 CLI 选项可能无法为该插件产生输出——因为它需要外部请求触发,进程常驻等待。

换句话说,要验证该插件,通常需要先启动 Telegraf,再向监听地址发起一次模拟的 Firehose 请求(下文"本地验证"小节会给出可直接使用的样例)。

[!IMPORTANT] 使用该插件会在 AWS 侧产生费用。CloudWatch Metric Streams 本身按流传输的指标数量计费(官方定价文档中的 "Metric Streams example" 一节有示例),启用前请评估成本。

二、整体工作流程

从源码 cloudwatch_metric_streams.go 的调用链可以清晰还原出数据流的处理管线:

  1. 启动监听Start()依据配置决定以普通 TCP 还是 TLS 方式在service_address上监听(约 L120-L152);注册插件时默认service_address = ":443"paths = ["/telegraf"](见init(),约 L419-L425)。
  2. 请求路由ServeHTTP先递增requests_received自监控计数,随后检查请求路径是否命中paths配置,未命中直接返回 404(约 L165-L177)。
  3. 可选鉴权authenticateIfSet在设置了access_key时校验请求头X-Amz-Firehose-Access-Key是否与配置值一致,不一致返回 401(约 L406-L417)。
  4. 解码与校验serveWrite(约 L211-L331)依次完成:请求体大小检查(超过max_body_size返回 413)→ 请求方法必须是 POST(否则 405)→ 若Content-Encoding: gzip则先解压 → JSON 反序列化请求结构(否则 400)。
  5. 逐条解析:遍历请求中的records,对每个record.data做 Base64 解码;解码结果是多行 JSON 拼接块(换行分隔),因此源码按\n切分后逐行json.Unmarshal成指标数据对象(约 L274-L306)。
  6. 组成指标composeMetrics将解码后的数据对象转换为 Telegraf metric(measurement、tags、fields、timestamp),写入 accumulator(约 L333-L371)。
  7. 回执:处理成功后,插件以请求中的requestId组装 JSON 响应并返回200 OK,让 Firehose 确认投递成功(约 L308-L330)。

三、完整配置与逐项解析

插件示例配置位于 sample.conf(README 中以@sample.conf内嵌),下面是完整配置及对应的源码行为说明:

# AWS Metric Streams listener [[inputs.cloudwatch_metric_streams]] ## Address and port to host HTTP listener on service_address = ":443" ## Paths to listen to. # paths = ["/telegraf"] ## maximum duration before timing out read of the request # read_timeout = "10s" ## maximum duration before timing out write of the response # write_timeout = "10s" ## Maximum allowed http request body size in bytes. ## 0 means to use the default of 524,288,000 bytes (500 mebibytes) # max_body_size = "500MB" ## Optional access key for Firehose security. # access_key = "test-key" ## An optional flag to keep Metric Streams metrics compatible with ## CloudWatch's API naming # api_compatability = false ## 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"

各配置项与源码的对应关系如下:

配置项对应源码字段默认值说明
service_addressServiceAddress":443"HTTP(S) 监听地址与端口,格式为host:port。若配置了 TLS 证书则按 HTTPS 服务启动
pathsPaths["/telegraf"]允许的监听路径列表。请求路径不在此列表中时返回 404。需要在 AWS 侧把 Firehose HTTP 端点指向这里
read_timeoutReadTimeout10s读取请求的超时时间。源码在Init()中会兜底:小于 1 秒时强制置为 10 秒(约 L106-L108)
write_timeoutWriteTimeout10s写回响应的超时时间,同样有"低于 1 秒则置为 10 秒"的兜底(约 L110-L112)
max_body_sizeMaxBodySize500MB(524,288,000 字节)允许的最大请求体字节数。0 表示使用默认值(常量defaultMaxBodySize,见约 L30-L33);超出后返回 HTTP 413
access_keyAccessKey空(不鉴权)可选的 Firehose 访问密钥。设置后要求请求头X-Amz-Firehose-Access-Key与之完全一致,否则返回 401
api_compatabilityAPICompatabilityfalse见下文"API 兼容模式";开启后统计字段会被重命名为 CloudWatch API 的命名
tls_allowed_cacertsServerConfig.TLSAllowedCACerts允许的客户端 CA 证书列表。配置后启用双向 TLS(mTLS),客户端必须持有可验证的证书
tls_cert/tls_keyServerConfig.TLSCert/TLSKey服务端证书与私钥。二者配置后监听器以tls.Listen启动(见Start()约 L125-L133)

TLS 相关字段继承自plugins/common/tls中的ServerConfig(参见 server.go),与 Telegraf 其他 HTTP 类插件保持一致的服务端 TLS 语义。

3.1 认证与安全的两种叠加手段

  • Access Key 鉴权access_key是最轻量的防护。实现位于authenticateIfSet(约 L406-L417):读取X-Amz-Firehose-Access-Key请求头并做字符串比对。AWS Firehose HTTP 交付规范本身支持在投递请求中携带该头,与 Telegraf 的校验天然配套。
  • TLS / 双向 TLS:配置tls_certtls_key后,插件通过tls.Listen建立 HTTPS 监听;进一步配置tls_allowed_cacerts则要求客户端出示受信任的证书,实现"服务器与客户端互验"的强认证链路。

四、数据格式与指标组成规则

4.1 AWS 侧投递的原始格式

Firehose 发送给插件的 HTTP 请求体是 JSON,核心结构为requestIdtimestamprecords数组;其中每个record.dataBase64 编码的、多行 JSON 拼接的数据块(即每个 data 字段里可以包含多条以换行分隔的 JSON,每个请求里可以有多个 record)。仓库测试数据 testdata/record.json 提供了一个完整样例,其中data字段 Base64 解码后即为 README 中展示的单条指标:

{ "metric_stream_name": "sandbox-dev-cloudwatch-metric-stream", "account_id": "541737779709", "region": "us-west-2", "namespace": "AWS/EC2", "metric_name": "CPUUtilization", "dimensions": { "InstanceId": "i-0efc7ghy09c123428" }, "timestamp": 1651679580000, "value": { "max": 10.011666666666667, "min": 10.011666666666667, "sum": 10.011666666666667, "count": 1 }, "unit": "Percent" }

字段语义与解码逻辑(源码data结构,约 L66-L76;逐条解码逻辑约 L274-L306):

  • metric_stream_name:指标流名称(解码后写入数据结构但不会进入 metric,从composeMetrics看并未使用该字段);
  • account_id:AWS 账户 ID;
  • region:指标产生的地域;
  • namespace:指标命名空间,如AWS/EC2
  • metric_name:指标名,如CPUUtilization
  • dimensions:维度字典,键值均为字符串;
  • timestamp:Unix 毫秒时间戳(代码中除以 1000 得到秒后调用time.Unix);
  • value:该时间点上的统计聚合值,通常包含maxminsumcount
  • unit:计量单位(如PercentBytesCount),同样不直接进入 Telegraf 字段。

注意:由于 CloudWatch 自身处理链路存在延迟,解码后指标携带的时间戳通常比处理时刻早 3~5 分钟,这是正常现象,下游查询与告警需按指标自带时间戳而非接收时刻来对齐。

4.2 Tags(标签)生成规则

依据 README "Tags" 一节及源码composeMetrics(约 L363-L368):

  • dimensions字典中的所有键值对都会作为 tag 加入 metric,键名保持原样(例如InstanceId=i-0efc7ghy09c123428);
  • 固定追加两个 tag:accountId(取自account_id)与region(取自region)。

4.3 Measurements 与 Fields(测量名与字段)

  • measurement(测量名):由namespacemetric_name拼接而成:先把namespace中的/替换为_AWS/EC2AWS_EC2),再整体转为小写并与下划线连接(约 L338-L339)。因此AWS/EC2+CPUUtilization得到测量名aws_ec2_cpuutilization
  • fields(字段)value对象中的每个聚合键都成为字段,例如maxminsumcount
  • timestamp(时间戳):直接使用指标自带的毫秒时间戳。

4.4 API 兼容模式(api_compatability)

Metric Streams 原生输出的聚合字段名(max/min/count)与 CloudWatch API(GetMetricStatistics 等)返回的字段名(Maximum/Minimum/SampleCount)并不一致。若希望从 API 轮询平滑迁移到 Metric Streams、保持下游查询与仪表盘字段不变,可设置api_compatability = true

源码在composeMetrics中完成了这一重命名(约 L346-L361):

Metric Streams 字段API 兼容模式字段
maxmaximum
minminimum
countsamplecount
sumsum(保持不变)

五、输出示例

以 4.1 节的 JSON 为例,开启与关闭api_compatability时分别得到如下两种输出(README "Example Output" 原文):

标准 Metric Streams 格式(api_compatability = false):

aws_ec2_cpuutilization,accountId=541737779709,region=us-west-2,InstanceId=i-0efc7ghy09c123428 max=10.011666666666667,min=10.011666666666667,sum=10.011666666666667,count=1 1651679580000

API 兼容格式(api_compatability = true):

aws_ec2_cpuutilization,accountId=541737779709,region=us-west-2,InstanceId=i-0efc7ghy09c123428 maximum=10.011666666666667,minimum=10.011666666666667,sum=10.011666666666667,samplecount=1 1651679580000

注意:实测样例中dimensions里的InstanceId作为 tag 原样输出,而accountIdregion是插件固定附加的。

六、源码级行为验证:测试用例解读

仓库配套测试 cloudwatch_metric_streams_test.go 覆盖了插件的关键路径,可作为理解行为与回归验证的参考:

  • TestWriteHTTP(约 L168-L191):向监听器 POSTtestdata/record.json,期望返回 200,验证基本投递链路。
  • TestWriteHTTPSNoClientAuth/TestWriteHTTPSWithClientAuth(约 L37-L106):验证配置 TLS 证书后以 HTTPS 投递成功;进一步配置tls_allowed_cacerts后,只有携带受信客户端证书的请求才能成功(测试使用testutil的 PKI 体系,证书样例位于 testutil/pki)。
  • TestWriteHTTPSuccessfulAuth/TestWriteHTTPFailedAuth(约 L108-L166):验证access_key鉴权——请求头携带正确密钥返回 200,错误密钥返回 401。
  • TestWriteHTTPExactMaxBodySize/TestWriteHTTPVerySmallMaxBody(约 L218-L268):验证max_body_size边界——请求体恰好等于上限时成功,超过上限返回 413。
  • TestWriteHTTPGzippedData(约 L446-L471):通过设置Content-Encoding: gzip并 POSTtestdata/records.gz,验证 gzip 解压路径。
  • TestComposeMetrics/TestComposeAPICompatibleMetrics(约 L339-L443):直接构造data对象调用composeMetrics,断言关闭/开启api_compatability时生成的 measurement、tags、fields 与时间戳,与 README 的转换规则完全一致。
  • TestReceive404ForInvalidEndpoint(约 L270-L293):请求未配置的路径返回 404。
  • TestWriteHTTPInvalid/TestWriteHTTPEmpty(约 L295-L337):无法解析或空请求体返回 400。

这些测试揭示了两个值得一提的实现细节:请求体超限返回 413(http.StatusRequestEntityTooLarge)、非 POST 方法返回 405,且这些错误响应均为 JSON 格式并附带error字段;同时所有错误路径都会通过bad_requests自监控指标按状态码打点(见tooLarge/methodNotAllowed/badRequest,约 L373-L404)。

七、Troubleshooting:自监控指标与排障

该插件通过 Telegraf 的selfstat框架注册了内部指标(见Init(),约 L92-L100),带address标签(值为service_address),可用于观测接入健康度:

指标名含义
requests_received监听器收到的请求总数
writes_served成功服务的写入请求数(每次 POST 处理完成时递增)
bad_requests失败请求数,按错误状态码(status_code标签)区分,如 400/405/413
request_time单个请求处理耗时(纳秒累计)
age_max/age_min当前区间内指标时间戳与处理时刻的最大/最小年龄差。由于 CloudWatch 指标通常滞后 3~5 分钟,这两个指标可用于在时序管线中校正基于时间戳的延迟/时延测量

排障提示(README "Troubleshooting" 一节):

  • 每次收到无法解析的请求时,插件会记录具体错误日志,并向 AWS 返回非 200 状态码(400/405/413/401 等),Firehose 侧会据此标记投递失败。
  • 如果从 AWS 侧持续看到投递失败,可优先核对:监听路径是否与 Firehose HTTP 端点配置一致、max_body_size是否放不下大数据块、access_key是否与请求头匹配、TLS 证书链是否完整。
  • Firehose HTTP 交付自身的故障排查可参考 AWS Firehose 的 HTTP 交付排障文档(README 中给出了对应链接)。

八、本地验证与 AWS 对接要点

8.1 本地验证

由于插件是服务输入,建议先启动 Telegraf,再模拟 Firehose 请求验证。仓库 testdata/record.json 中的请求体即为合规样例,可将其 POST 到监听路径验证:

# 方式一:直接 POST 样例请求(默认监听 443 需要调整端口或改用非特权端口) # 先在 telegraf.conf 中把 service_address 改为如 ":8080",并保持 paths 默认值 curl -X POST -H 'Content-Type: application/json' \ -d @plugins/inputs/cloudwatch_metric_streams/testdata/record.json \ http://127.0.0.1:8080/telegraf # 期望返回 HTTP 200,并回显 {"requestId":"c8291d2e-8c46-4f2a-a8df-2562550287ad","timestamp":<毫秒时间戳>}

若在配置中开启了access_key,则请求头需追加-H 'X-Amz-Firehose-Access-Key: <你的密钥>';若开启了 TLS,则改用https://并携带相应证书。

8.2 AWS 侧对接要点

  1. 在 CloudWatch 控制台创建Metric Stream,选择投递目标为Kinesis Data Firehose,输出格式选择"OpenTelemetry 或 JSON"中的 JSON 透传(插件按 Firehose HTTP 交付规范解析)。
  2. 配置 Firehose 的HTTP endpoint指向 Telegraf 的service_address+paths(默认https://<telegraf主机>:443/telegraf),并可设置访问密钥(与access_key对应)。
  3. 确认安全组/防火墙放行对应端口,Telegraf 主机需能被 AWS 侧访问。
  4. 观察requests_receivedwrites_served是否持续增长,并用age_max评估端到端延迟。

8.3 其他集成注意点

  • 插件注册文件位于 plugins/inputs/all/cloudwatch_metric_streams.go,构建标签为inputs.cloudwatch_metric_streams,使用 Telegraf 自定义构建(custom builder)时可按需裁剪。
  • 插件全局配置(如name_overridetagsfieldpass等通用选项)同样适用于该插件,详见 docs/CONFIGURATION.md 中关于插件的通用配置章节。

九、总结

cloudwatch_metric_streams是一个"以 HTTP 服务端姿态对接 AWS 指标推送"的服务输入插件:它把 Firehose HTTP 交付的 Base64+JSON 数据流,安全(可选的 Access Key 与双向 TLS)、可靠(413/405/400 分级错误响应与自监控指标)地转换为 Telegraf 标准指标。其核心使用要点可归纳为四条:

  1. 路径与端口要对齐service_address+paths必须与 Firehose HTTP 端点严格一致;
  2. 字段命名有两个模式:默认保留 Metric Streams 原生聚合名,api_compatability = true时切换为 CloudWatch API 命名(maximum/minimum/samplecount),平滑迁移场景优先开启;
  3. 指标自带时间戳:比处理时刻滞后 3~5 分钟属正常现象,延迟观测请使用age_max/age_min自监控指标;
  4. 本地调试有现成样例:直接复用 testdata/record.json 与 cloudwatch_metric_streams_test.go 即可快速验证接入链路。

【免费下载链接】telegrafAgent for collecting, processing, aggregating, and writing metrics, logs, and other arbitrary data.项目地址: https://gitcode.com/GitHub_Trending/te/telegraf

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/15 1:42:55

Rust 1.XX (beta) -> Rust 1.XX

Rust 1.XX (beta) -> ## Rust 1.XX 【免费下载链接】rust-clippy A bunch of lints to catch common mistakes and improve your Rust code. Book: https://doc.rust-lang.org/clippy/ 项目地址: https://gitcode.com/GitHub_Trending/ru/rust-clippy - 更新新稳定版本…

作者头像 李华
网站建设 2026/9/15 1:41:20

RAG工程落地全链路实战:从文档切块到K8s生产部署

1. 项目概述&#xff1a;这不是“速成课”&#xff0c;而是一份RAG工程落地的完整施工图你点开这个标题&#xff0c;第一反应可能是——又一个标题党&#xff1f;7天从小白到大神&#xff1f;吊打付费&#xff1f;存下吧很难找全&#xff1f;这些话术确实刺眼&#xff0c;但如果…

作者头像 李华
网站建设 2026/9/15 1:40:11

用Doom实测Astra云电脑:老游戏才是串流延迟的照妖镜

现在测云电脑的人&#xff0c;第一反应都是打开 3A 大作&#xff0c;画面一糊、帧率一掉就断定平台不行。但我一直觉得这个思路反了——真正能看出一个串流平台底子的&#xff0c;恰恰是那些“看起来毫无压力”的老游戏。Doom 这种老祖宗级别的 FPS&#xff0c;对帧率天花板要求…

作者头像 李华
网站建设 2026/9/15 1:39:21

Flutter与OpenHarmony构建高性能播放器进度条实践

1. 为什么选择 Flutter OpenHarmony 构建播放器控件在移动端开发领域&#xff0c;播放器进度条看似简单&#xff0c;实则涉及跨平台渲染性能、手势交互精度、状态同步等复杂问题。传统方案通常面临三个困境&#xff1a;一是原生开发需要针对Android/iOS分别实现&#xff0c;维…

作者头像 李华
网站建设 2026/9/15 1:37:26

实体关系抽取实战:从依赖树到图卷积神经网络的完整实现

简介&#xff1a;基于图卷积神经网络的实体关系抽取项目&#xff0c;面向深度学习、自然语言处理方向的在校学生、研究人员及企业开发者&#xff0c;完整覆盖实体关系抽取中数据预处理、GCN模型构建、训练测试、结果评估与可视化展示的流程。整个资源包共41个文件&#xff0c;以…

作者头像 李华