
Vector 的 AWS 数据出口方案aws_kinesis_firehose、aws_s3、aws_sqs 与 aws_ecs_metrics 五大集成深度解析【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector将可观测性数据搬出 AWS 往往需要把 CloudWatch Logs 订阅、Kinesis Firehose、S3 备份、SQS 等多个服务拼接成一条脆弱的数据管道。Vector 在 0.11.0 版本中一次性交付了五个互补的 AWS 集成组件aws_kinesis_firehose数据源、aws_cloudwatch_logs_subscription_parser转换该组件已在 v0.23.0 中移除、aws_s3数据源、aws_sqs数据汇以及aws_ecs_metrics数据源见发布亮点 website/content/en/highlights/2020-10-28-new-aws-integrations.md。本文以该亮点文档为主干结合当前仓库中这些组件的源码实现与配套指南逐组件讲清配置项、默认值与底层处理机制并给出一条CloudWatch Logs → Kinesis Firehose → Vector → 任意数据汇的可复制生产级链路。五大集成分工总览组件类型作用源码位置aws_kinesis_firehosesource以 HTTPS 端点接收 Kinesis Firehose 投递的 HTTP 事件src/sources/aws_kinesis_firehose/mod.rsaws_cloudwatch_logs_subscription_parsertransform解析 CloudWatch Logs 订阅事件v0.23.0 起移除由 VRL 函数替代—aws_s3source消费写入 S3 的日志对象经 SQS 通知驱动src/sources/aws_s3/mod.rsaws_sqssink将日志/事件发布到 Amazon SQS 队列src/sinks/aws_s_s/sqs/config.rsaws_ecs_metricssource周期性抓取 ECS/Fargate 任务元数据端点输出容器指标src/sources/aws_ecs_metrics/mod.rs这五个组件覆盖了数据出 AWS的三个典型场景Firehose 主动推送aws_kinesis_firehose、对象存储落盘后拉取aws_s3aws_sqs、以及容器指标采集aws_ecs_metrics。下面逐一深入。aws_kinesis_firehose 数据源把 Firehose 变成 Vector 的 HTTP 入口aws_kinesis_firehose数据源在 Vector 内部实现为一个 HTTPS 服务器从源码结构看它在 build 逻辑 中基于 hyper 构建Server通过warp挂载 Firehose 专用的请求过滤器filters.rs支持 TLS 终结、keepalive 与优雅停机。Kinesis Firehose 的 HTTP Endpoint Destination 只会向 443 端口的 HTTPS 地址投递数据因此通常需要把 Vector 实例放在负载均衡器后面或直接在该数据源上配置tls参数提供有效证书。核心配置项结合 AwsKinesisFirehoseConfig 定义主要配置项如下address必填监听地址如0.0.0.0:443或localhost:443access_keysFirehose 可随每次请求携带用户自定义 access key通过x-amz-firehose-access-key请求头。配置后只接受匹配 key 的请求不配置则放行所有请求。旧字段access_key单数已标记为 deprecated源码中会自动将其合并进access_keys并打印弃用警告store_access_key为true时请求中携带的 access key 会作为aws_kinesis_firehose_access_key存入事件 secretsrecord_compression取值为auto默认/none/gzip控制对 Firehose 消息内每条记录的解压方式。注意这与 Firehose 端 HTTP 请求整体的 Content-Encoding 是两个概念CloudWatch Logs 发出的记录会先做 gzip 压缩再 base64 编码而整包请求体是否 gzip 由RecordConfiguration.ContentEncoding决定common_attributes从X-Amz-Firehose-Common-Attributes请求头中挑选属性注入事件支持通配符如environment、application_*、*framing/decoding标准的分帧与解码配置决定记录内容最终如何变成日志事件acknowledgements该数据源返回can_acknowledge true支持端到端确认keepalive、log_namespaceTCP keepalive 与日志命名空间Vector 命名空间 / Legacy 命名空间控制。事件输出形态从 outputs 方法 与同文件内的单元测试如aws_kinesis_firehose_forwards_events_vector_namespace可以确认每条事件携带source_type aws_kinesis_firehose、解码后的messagegzip base64 记录被还原后的原文以及request_id、source_arn来自x-amz-firehose-request-id、x-amz-firehose-source-arn请求头等源元数据在 Vector 命名空间下这些字段挂在aws_kinesis_firehose元数据路径下common_attributes匹配结果也放在同一前缀下。测试同时验证了压缩策略的边界行为record_compression: gzip时收到未压缩记录会返回 400而auto按 magic bytes 探测后失败则原样转发。CloudWatch Logs 经 Firehose 流入 Vector 的完整链路仓库中配套指南 website/content/en/guides/aws/cloudwatch-logs-firehose.md 给出了生产可用链路CloudWatch Logs 订阅过滤器 → Kinesis FirehoseHTTP Endpoint Destination→ Vector → 任意数据汇。该链路的优势在于Firehose 会按可配置时长重试并把失败事件死信到 S3、多个 Vector 实例可放在负载均衡器后分摊流量、一条交付流可承接多个日志组的订阅。Vector 侧配置sources: firehose: type: aws_kinesis_firehose address: 0.0.0.0:8080 # 公网地址在配置 Firehose 时指定 access_key: ${FIREHOSE_ACCESS_KEY} # 与 Firehose 配置中保持一致 transforms: parse: type: remap inputs: [firehose] drop_on_error: false source: | parsed parse_aws_cloudwatch_log_subscription_message!(.message) . unnest(parsed.log_events) . map_values(.) - |value| { event del(value.log_events) value | event message del(.message) . | object!(parse_json!(message)) } sinks: console: type: console inputs: [parse] encoding: codec: json这里值得强调一点原 0.11.0 高亮文档中的aws_cloudwatch_logs_subscription_parser转换组件负责的工作在当前版本中已由 VRL 内置函数parse_aws_cloudwatch_log_subscription_message承担该转换已在 v0.23.0 移除。上面这段 remap 程序先解析.message得到订阅事件对象再用unnest(parsed.log_events)把一条订阅事件拆成多条独立日志事件最后用map_values迭代把嵌套的id/timestamp提升到根层并将message字段按 JSON 展开。AWS 侧搭建步骤指南中的关键 CLI 步骤省略环境变量说明部分依次为创建日志组与 Firehose 调试日志组aws logs create-log-group --log-group-name ${LOG_GROUP} aws logs create-log-group --log-group-name ${FIREHOSE_LOG_GROUP} aws logs create-log-stream \ --log-group-name ${FIREHOSE_LOG_GROUP} \ --log-stream-name ${FIREHOSE_LOG_STREAM}创建用于存放失败事件的 S3 桶HTTP Endpoint Destination 强制要求aws s3api create-bucket --bucket ${FIREHOSE_S3_BUCKET} \ $(if [[ ${AWS_REGION} ! us-east-1 ]] ; then echo --create-bucket-configuration LocationConstraint${AWS_REGION} ; fi)创建允许 Firehose 服务的 IAM 角色FirehoseVector并挂载允许写 S3 桶和调试日志组的策略创建指向 Vector 的交付流核心参数包括EndpointConfiguration.UrlVector 公网地址、AccessKey与 Vector 配置一致、RetryOptions.DurationInSeconds: 300、S3BackupMode: FailedDataOnly300 秒内投递失败的事件落 S3aws firehose create-delivery-stream --delivery-stream-name ${FIREHOSE_DELIVERY_STREAM} \ --http-endpoint-destination file://(cat EOF { EndpointConfiguration: { Url: ${VECTOR_ENDPOINT}, Name: vector, AccessKey: ${FIREHOSE_ACCESS_KEY} }, RequestConfiguration: { ContentEncoding: GZIP }, CloudWatchLoggingOptions: { Enabled: true, LogGroupName: ${FIREHOSE_LOG_GROUP}, LogStreamName: ${FIREHOSE_LOG_STREAM} }, RoleARN: arn:aws:iam::${AWS_ACCOUNT_ID}:role/FirehoseVector, RetryOptions: { DurationInSeconds: 300 }, S3BackupMode: FailedDataOnly, S3Configuration: { RoleARN: arn:aws:iam::${AWS_ACCOUNT_ID}:role/FirehoseVector, BucketARN: arn:aws:s3:::${FIREHOSE_S3_BUCKET} } } EOF )为 CloudWatch Logs 创建订阅角色CWLtoKinesisFirehoseRole允许firehose:*指向该交付流再创建订阅过滤器aws logs put-subscription-filter \ --log-group-name ${LOG_GROUP} \ --filter-name Destination \ --filter-pattern \ --destination-arn arn:aws:firehose:${AWS_REGION}:${AWS_ACCOUNT_ID}:deliverystream/${FIREHOSE_DELIVERY_STREAM} \ --role-arn arn:aws:iam::${AWS_ACCOUNT_ID}:role/CWLtoKinesisFirehoseRole指南还附带了验证方法用另一份 Vector 配置stdin数据源 aws_cloudwatch_logs数据汇向日志组灌入数据或在 CloudWatch 控制台中手动创建 log event观察事件在约 300 秒的批处理间隔后出现在 Vector 侧。指南的 Deep dive 章节演示了同一批 JSON 事件在 Firehose 请求体、aws_kinesis_firehose解码后事件、remap 拆解后事件三个阶段的形态演变适合用来理解字段在管道中如何逐步展开。aws_s3 数据源SQS 通知驱动的对象消费aws_s3数据源src/sources/aws_s3/mod.rs用于消费被写入 S3 的日志对象例如 Firehose 失败事件备份桶、其他组件直接写入 S3 的日志文件。从源码结构看它目前只实现了一种Strategy::Sqs策略Strategy 枚举监听一个 SQS 队列中的 S3 桶通知事件收到通知后去 S3 拉取对应对象并解码——SQS 子配置位于 src/sources/aws_s3/sqs.rs。关键配置项结合 AwsS3Config 定义region区域可通过RegionOrEndpoint同时指定endpoint指向兼容服务compressionauto默认/none/gzip/zstd。auto模式按顺序探测Content-Encoding元数据、Content-Type元数据、对象 key 的文件扩展名.gz/.zst三者都识别不出则按未压缩处理——该探测逻辑见 determine_compression并有同名单元测试覆盖sqsSQS 队列消费配置队列地址、轮询间隔、批量大小、可见性超时、delete_failed_message等未配置时构建会报strategysqs 需要 sqs 配置错误multiline可选的多行聚合配置request_payer设为requester表示 Vector 所在 AWS 账号承担请求者付费桶Requester Pays的请求与流量费用force_path_style默认true控制桶名放在主机名还是 URL 路径中authAWS 鉴权配置旧的assume_role字段已 deprecatedacknowledgements支持端到端确认事件处理失败时可保留 SQS 消息不删除以便重试。从 outputs 方法 可见事件会附带bucket、object、region、timestamp等源元数据便于下游按对象溯源。aws_sqs 数据汇向 SQS 队列发布事件aws_sqs数据汇src/sinks/aws_s_s/sqs/config.rs与aws_sns数据汇共享 BaseSSSinkConfig 基础配置。核心字段queue_url必填URI 格式目标 SQS 队列地址如https://sqs.us-east-2.amazonaws.com/123456789012/MyQueueregion队列所在区域encoding事件编码codec输入类型限制为日志message_group_id仅适用于 FIFO 队列的模板。校验规则在 message_group_id 函数 中实现队列 URL 以.fifo结尾时必须提供该字段否则报message_group_id should be defined for FIFO queue非 FIFO 队列提供则会报is not allowed with non-FIFO queue同文件中的测试validate_rejects_fifo_without_message_group_id验证了这一行为message_deduplication_idFIFO 队列消息去重 ID 模板需为每条事件生成唯一字符串auth/tls/request/acknowledgements标准的 AWS 鉴权、TLS、请求重试与确认配置。健康检查实现为对目标队列调用get_queue_attributeshealthcheck 函数在 Vector 启动或vector validate时验证队列可达且凭据有效。aws_ecs_metrics 数据源ECS/Fargate 容器指标aws_ecs_metrics数据源src/sources/aws_ecs_metrics/mod.rs用于采集 AWS ECS 与 AWS Fargate 任务的 Docker 容器指标配置项与行为如下endpoint任务元数据端点基地址。留空时自动探测——文档语义上优先使用环境变量ECS_CONTAINER_METADATA_URI_V4V4 端点其次ECS_CONTAINER_METADATA_URIV3都未设置时回退到 V2 的固定地址169.254.170.2/v2见 default_endpoint 实现versionv2/v3/v4为空时按上述环境变量规则自动推断scrape_interval_secs抓取间隔默认15 秒namespace指标命名空间默认awsecs设为空字符串则禁用命名空间。从源码看统计端点的拼接规则是V2 版本请求{endpoint}/statsV3/V4 版本请求{endpoint}/task/statsstats_endpoint 方法。抓取循环aws_ecs_metrics 函数用定时间隔流驱动 HTTP GET响应体交给 parser.rs 解析为指标事件并批量发出解析失败、非 200 响应或网络错误分别发出AwsEcsMetricsParseError、HttpClientHttpResponseError、HttpClientHttpError内部事件。该数据源输出指标类型SourceOutput::new_metrics()且不支持确认can_acknowledge false。sources: ecs_metrics: type: aws_ecs_metrics # endpoint/version 可省略由容器内环境变量自动探测 scrape_interval_secs: 15 namespace: awsecs小结0.11.0 引入的这套 AWS 集成本质上是围绕三条路径解决数据出 AWSFirehose 推送适合低延迟、多日志组汇聚的实时链路配合指南中的 S3 死信与重试机制保证韧性S3 经 SQS 通知的消费适合对象化落盘后的批量摄取aws_sqs数据汇则把 Vector 处理结果送进 AWS 消息体系供其他系统消费aws_ecs_metrics补齐了容器指标的采集。需要留意版本差异0.11.0 时代的aws_cloudwatch_logs_subscription_parser转换已移除当前版本请统一使用 VRL 的parse_aws_cloudwatch_log_subscription_message函数完成订阅事件解析各组件的access_key/assume_role等旧字段也已在源码中标记 deprecated新配置应使用access_keys与auth。所有上述行为均可在仓库对应源码文件与内置测试中直接核对。【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考