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,其关键差异有两点:
- 全局或插件级
interval设置可能不生效——数据的到达节奏由上游 Firehose 推送频率决定; --test、--test-wait、--once等 CLI 选项可能无法为该插件产生输出——因为它需要外部请求触发,进程常驻等待。
换句话说,要验证该插件,通常需要先启动 Telegraf,再向监听地址发起一次模拟的 Firehose 请求(下文"本地验证"小节会给出可直接使用的样例)。
[!IMPORTANT] 使用该插件会在 AWS 侧产生费用。CloudWatch Metric Streams 本身按流传输的指标数量计费(官方定价文档中的 "Metric Streams example" 一节有示例),启用前请评估成本。
二、整体工作流程
从源码 cloudwatch_metric_streams.go 的调用链可以清晰还原出数据流的处理管线:
- 启动监听:
Start()依据配置决定以普通 TCP 还是 TLS 方式在service_address上监听(约 L120-L152);注册插件时默认service_address = ":443"、paths = ["/telegraf"](见init(),约 L419-L425)。 - 请求路由:
ServeHTTP先递增requests_received自监控计数,随后检查请求路径是否命中paths配置,未命中直接返回 404(约 L165-L177)。 - 可选鉴权:
authenticateIfSet在设置了access_key时校验请求头X-Amz-Firehose-Access-Key是否与配置值一致,不一致返回 401(约 L406-L417)。 - 解码与校验:
serveWrite(约 L211-L331)依次完成:请求体大小检查(超过max_body_size返回 413)→ 请求方法必须是 POST(否则 405)→ 若Content-Encoding: gzip则先解压 → JSON 反序列化请求结构(否则 400)。 - 逐条解析:遍历请求中的
records,对每个record.data做 Base64 解码;解码结果是多行 JSON 拼接块(换行分隔),因此源码按\n切分后逐行json.Unmarshal成指标数据对象(约 L274-L306)。 - 组成指标:
composeMetrics将解码后的数据对象转换为 Telegraf metric(measurement、tags、fields、timestamp),写入 accumulator(约 L333-L371)。 - 回执:处理成功后,插件以请求中的
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_address | ServiceAddress | ":443" | HTTP(S) 监听地址与端口,格式为host:port。若配置了 TLS 证书则按 HTTPS 服务启动 |
paths | Paths | ["/telegraf"] | 允许的监听路径列表。请求路径不在此列表中时返回 404。需要在 AWS 侧把 Firehose HTTP 端点指向这里 |
read_timeout | ReadTimeout | 10s | 读取请求的超时时间。源码在Init()中会兜底:小于 1 秒时强制置为 10 秒(约 L106-L108) |
write_timeout | WriteTimeout | 10s | 写回响应的超时时间,同样有"低于 1 秒则置为 10 秒"的兜底(约 L110-L112) |
max_body_size | MaxBodySize | 500MB(524,288,000 字节) | 允许的最大请求体字节数。0 表示使用默认值(常量defaultMaxBodySize,见约 L30-L33);超出后返回 HTTP 413 |
access_key | AccessKey | 空(不鉴权) | 可选的 Firehose 访问密钥。设置后要求请求头X-Amz-Firehose-Access-Key与之完全一致,否则返回 401 |
api_compatability | APICompatability | false | 见下文"API 兼容模式";开启后统计字段会被重命名为 CloudWatch API 的命名 |
tls_allowed_cacerts | ServerConfig.TLSAllowedCACerts | 空 | 允许的客户端 CA 证书列表。配置后启用双向 TLS(mTLS),客户端必须持有可验证的证书 |
tls_cert/tls_key | ServerConfig.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_cert与tls_key后,插件通过tls.Listen建立 HTTPS 监听;进一步配置tls_allowed_cacerts则要求客户端出示受信任的证书,实现"服务器与客户端互验"的强认证链路。
四、数据格式与指标组成规则
4.1 AWS 侧投递的原始格式
Firehose 发送给插件的 HTTP 请求体是 JSON,核心结构为requestId、timestamp与records数组;其中每个record.data是Base64 编码的、多行 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:该时间点上的统计聚合值,通常包含max、min、sum、count;unit:计量单位(如Percent、Bytes、Count),同样不直接进入 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(测量名):由
namespace与metric_name拼接而成:先把namespace中的/替换为_(AWS/EC2→AWS_EC2),再整体转为小写并与下划线连接(约 L338-L339)。因此AWS/EC2+CPUUtilization得到测量名aws_ec2_cpuutilization。 - fields(字段):
value对象中的每个聚合键都成为字段,例如max、min、sum、count。 - 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 兼容模式字段 |
|---|---|
max | maximum |
min | minimum |
count | samplecount |
sum | sum(保持不变) |
五、输出示例
以 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 1651679580000API 兼容格式(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 原样输出,而accountId、region是插件固定附加的。
六、源码级行为验证:测试用例解读
仓库配套测试 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 侧对接要点
- 在 CloudWatch 控制台创建Metric Stream,选择投递目标为Kinesis Data Firehose,输出格式选择"OpenTelemetry 或 JSON"中的 JSON 透传(插件按 Firehose HTTP 交付规范解析)。
- 配置 Firehose 的HTTP endpoint指向 Telegraf 的
service_address+paths(默认https://<telegraf主机>:443/telegraf),并可设置访问密钥(与access_key对应)。 - 确认安全组/防火墙放行对应端口,Telegraf 主机需能被 AWS 侧访问。
- 观察
requests_received与writes_served是否持续增长,并用age_max评估端到端延迟。
8.3 其他集成注意点
- 插件注册文件位于 plugins/inputs/all/cloudwatch_metric_streams.go,构建标签为
inputs.cloudwatch_metric_streams,使用 Telegraf 自定义构建(custom builder)时可按需裁剪。 - 插件全局配置(如
name_override、tags、fieldpass等通用选项)同样适用于该插件,详见 docs/CONFIGURATION.md 中关于插件的通用配置章节。
九、总结
cloudwatch_metric_streams是一个"以 HTTP 服务端姿态对接 AWS 指标推送"的服务输入插件:它把 Firehose HTTP 交付的 Base64+JSON 数据流,安全(可选的 Access Key 与双向 TLS)、可靠(413/405/400 分级错误响应与自监控指标)地转换为 Telegraf 标准指标。其核心使用要点可归纳为四条:
- 路径与端口要对齐:
service_address+paths必须与 Firehose HTTP 端点严格一致; - 字段命名有两个模式:默认保留 Metric Streams 原生聚合名,
api_compatability = true时切换为 CloudWatch API 命名(maximum/minimum/samplecount),平滑迁移场景优先开启; - 指标自带时间戳:比处理时刻滞后 3~5 分钟属正常现象,延迟观测请使用
age_max/age_min自监控指标; - 本地调试有现成样例:直接复用 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),仅供参考