Vector aggregate 转换组件详解:指标聚合模式、配置参数与事件时间聚合机制
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
aggregate是 Vector 拓扑中的一个有状态指标转换组件,用于在一个可配置的间隔窗口内,把同一指标序列(name、namespace、tags 等)的多个 metric 事件合并为单个事件,从而以降低时间粒度的代价换取指标数据量的显著缩减。它支持全部六类指标(counter、gauge、set、summary、histogram、distribution),但不接受 logs 与 traces 输入。掌握本组件后,你可以在高吞吐链路中减少下游 sink 的处理与传输开销,并在指标时间戳敏感的接收端(如 Datadog Metrics sink)通过事件时间聚合(event-time aggregation)避免不同样本被合并成同一个时间点。
组件定位与输入输出
从组件元数据 aggregate.cue 可以看到 aggregate 的关键属性:
egress_method: "stream":作为流式转换组件运行,事件逐条经过、按窗口批量吐出;stateful: true:组件内部维护聚合状态(每个序列当前累积值),因此它是有状态的;- 输入支持
metric的全部子类型:counter、distribution、gauge、histogram、set、summary,logs 与 traces 均为false; - 输出为单一
""输出端口,内容是修改后的输入metric事件。
组件配置结构定义在 AggregateConfig,并注册为名为aggregate的转换组件(impl TransformConfig位于 config.rs)。
配置参数说明
aggregate 的完整配置面由三个顶层项组成:interval_ms、mode、event_time。下面逐项说明,取值与默认值均与 config.rs 中的源码一致。
interval_ms:刷新间隔
- 类型:
uint(毫秒),默认10000(10 秒),见 default_interval_ms; - 必须大于 0,且不能超过
i64::MAX毫秒——Aggregate::new 在构造时就会对interval_ms、event_time.max_future_ms、event_time.allowed_lateness_ms做范围校验,超限直接报错拒绝构建组件; - 语义:每次 flush 之间的时间间隔。在这个时间框架内,具有相同序列数据(name、namespace、tags 等)的指标事件会被聚合到一起。
mode:聚合函数
- 类型:字符串枚举,默认
Auto,定义于 AggregationMode; - 各取值含义:
| mode | 行为 |
|---|---|
Auto(默认) | 增量(incremental)指标求和,绝对(absolute)指标取最新值 |
Sum | 只聚合增量指标(求和);绝对指标原样透传 |
Latest | 绝对指标取最新值;增量指标原样透传 |
Count | 对增量与绝对指标都计数(样本条数) |
Diff | 对绝对指标返回最新值与上一窗口的差值;增量指标原样透传 |
Max/Min | 绝对指标取窗口内最大值/最小值;增量指标原样透传 |
Mean/Stdev | 对绝对指标计算均值/标准差(仅处理 gauge 值);增量指标原样透传 |
注意“原样透传”(pass through unchanged)这一行为:例如在Sum模式下收到的absolute指标不会被存储或聚合,而是立即转发到下游。源码中这一判定集中在 passes_through_unchanged;另外Mean/Stdev模式下的绝对非 gauge 值既不聚合也不透传,而是被静默忽略,见 is_silently_ignored。
event_time:事件时间聚合块(可选)
出现该配置块即启用事件时间聚合:指标按每个事件自身的时间戳分桶,而不是按 Vector 处理它们的时刻分桶。省略该块则保持默认的系统时间(system-time)行为。子项定义在 EventTimeConfig:
allowed_lateness_ms(uint,默认0):迟到事件的宽限期。每个桶会一直接收事件,直到系统时钟到达bucket_end + allowed_lateness_ms(bucket_end为该事件时间窗口的不含端点终点)。这个截止是在记录事件时强制的,而不只是周期性 flush 运行时;桶一旦发出就永久关闭,之后再落入该桶的事件会被丢弃并通过component_discarded_events_total计数。设为0表示严格序(不允许任何迟到)。missing_timestamp(枚举,默认drop):处理缺失时间戳的事件。需要被分桶的指标,取drop时直接丢弃并递增component_discarded_events_total;取use_system_time时用当前系统时钟合成时间戳。透传指标(见上文 mode 规则)不要求有时间戳。max_future_ms(uint,默认10000毫秒,见 default_max_future_ms):时钟漂移保护。时间戳领先系统时钟超过该值的事件会被丢弃并计数。设为0表示允许任意未来时间。
完整示例:5 秒窗口聚合
下面的示例来自组件元数据 aggregate.cue,配置为interval_ms: 5000(未指定mode即默认Auto):
输入(5 条事件):
# 输入 1:counter.1 @ 07:58:44.223543Z,host=my.host.com,incremental,value=1.1 # 输入 2:counter.1 @ 07:58:45.223543Z,host=my.host.com,incremental,value=2.2 # 输入 3:counter.1 @ 07:58:45.223543Z,host=different.host.com,incremental,value=1.1 # 输入 4:gauge.1 @ 07:58:47.223543Z,host=my.host.com,absolute,value=22.33 # 输入 5:gauge.1 @ 07:58:45.223543Z,host=my.host.com,absolute,value=44.55配置:
transforms: agg: type: aggregate interval_ms: 5000输出(3 条事件):
# counter.1 @ 07:58:45.223543Z,host=my.host.com,incremental,value=3.3(1.1 + 2.2 求和) # counter.1 @ 07:58:45.223543Z,host=different.host.com,incremental,value=1.1(独立序列,不受影响) # gauge.1 @ 07:58:45.223543Z,host=my.host.com,absolute,value=44.55(时间戳较晚的 47 秒那条覆盖了 45 秒的 22.33)这个示例直观展示了Auto模式的三条规则:同序列的增量值求和、tags 不同的序列各自独立、同序列的绝对值取“最新”(系统时间语义下即最后到达的那条)。
聚合行为详解与源码印证
组件文档描述的聚合规则(incremental指标在窗口内“相加”,absolute指标用新值替换旧值,保持数值正确性)在 transform.rs 中逐条落地:
- 状态容器:
Aggregate结构体持有map: HashMap<MetricSeries, MetricEntry>用于系统时间模式下的按序列聚合,见 Aggregate 结构定义; - 记录路径:record() 按
(mode, kind)匹配分发——Auto/Sum下的增量指标进入 record_sum()(要求新旧kind相同才可累加,否则发出AggregateUpdateFailed并以新值覆盖);Count模式调用 record_count(),对每条样本累加 1;Max/Min调用 record_comparison() 仅对 gauge 值做比较;Mean/Stdev把绝对 gauge 样本暂存到multi_map,在 flush 时求均值或标准差(见 flush_system_time() 中value /= entries.len()与方差开方逻辑); - 刷新循环:TaskTransform 实现 用
tokio::time::interval驱动定时 flush,tokio::select!在“到点刷新”与“收到新事件”之间交替;当输入流关闭(shutdown 或拓扑 reload)时调用 flush_final() 把状态中所有剩余指标吐出,保证在途指标不会被静默丢弃。
文档中给出的数值例子——两条 10 和 13 的incrementalcounter 聚合为 23,两条 93 和 95 的absolutegauge 取 95——正是record_sum累加与map.insert替换两处逻辑的直接结果。聚合的收益则在于数据量缩减:在按指标量计费、受处理 CPU 或网络带宽约束的场景中,它能直接降低成本,也可以减轻 aggregate 下游 transform 与 sink 的处理压力。
事件时间聚合机制
当配置了event_time块时,处理路径切换到 record_event_time(),核心机制如下:
分桶对齐
桶边界按 Unix 纪元起interval_ms的整数倍对齐,由 bucket_key() 计算:timestamp_ms.div_euclid(interval_ms).saturating_mul(interval_ms)。使用欧几里得除法是为了正确处理纪元之前(负毫秒)的时间戳。同一来源时间戳无论 Vector 何时收到,总是映射到同一个桶。
水位线与迟到事件
组件维护一个watermark,定义为“最近已发出桶的不含端点终点”,等价于仍允许写入的最小桶键。判断逻辑在 was_bucket_flushed():只要bucket_key < watermark,事件即被丢弃。而allowed_lateness_ms由 is_past_bucket_cutoff() 在记录时执行:当系统时钟到达bucket_end + allowed_lateness_ms之后,即使周期性 flush 尚未运行,新事件也会被拒绝——因此allowed_lateness_ms = 0能在 flush 间隔较长或不对齐时仍强制严格序。水位线在 flush_event_time_buckets() 中随每次 flush 推进到已发出桶的终点。
缺失与未来时间戳
- 需要分桶的指标缺失时间戳时默认丢弃;
use_system_time则用系统时钟合成。为保持判定一致,record_event_time() 只取一次Utc::now(),同时用于合成时间戳与后续漂移/截止比较; - 时间戳超前系统时钟超过
max_future_ms的事件作为时钟偏移保护被丢弃(max_future_ms = 0关闭该保护,此时 flush 端使用饱和加法防止i64回绕,见 flush_event_time_buckets 注释); - 所有丢弃都通过
component_discarded_events_total计数(丢弃原因记录在日志中,不打在指标标签上)。
透传与忽略规则的一致性
事件时间模式与系统时间模式保持一致的透传/忽略规则:mode不聚合的指标(如mean模式下的incremental事件、sum模式下的absolute事件)原样透传,不创建桶、不影响水位线;mean/stdev模式下的绝对非 gauge 值被忽略。这一行为由 passes_through_unchanged() 与 will_be_stored() 在创建桶之前判定——这很重要,因为空桶也会被视为可 flush 并推进水位线,若为不可存的事件无谓建桶,可能误拒之前桶内的合法事件。
另一个事件时间语义差异:Auto/Latest/Diff模式下的“latest”指桶内事件时间戳最新的样本(见 select_latest_by_event_timestamp()),而非最后到达的样本。
关闭与重载
输入流关闭时(shutdown 或拓扑重载),所有剩余事件时间桶都会先 flush 再退出。Diff模式额外保留一小段滚动窗口的前置桶(event_time_prev_buckets 仅存MetricData、丢弃 metadata 以避免持有 ack 终结器),用于计算跨桶边界的差值;其他模式不保留前置桶。
典型使用场景
事件时间聚合解决的是“下游 sink 以指标时间戳为键”时的数据坍缩问题。典型例子是 Datadog Metrics sink 对相同时间戳会覆盖旧值:按事件时间分桶可以防止不同样本折叠为同一个点。该特性在变更日志 24410_aggregate_event_time_aggregation.feature.md 中亦有说明。
内部遥测指标
aggregate 组件通过 内部事件 上报以下遥测:
| 指标 | 含义 |
|---|---|
aggregate_events_recorded_total | 被组件记录(进入聚合状态)的事件总数 |
aggregate_flushes_total | 完成一次 flush 的次数 |
aggregate_failed_updates | 聚合更新失败次数(如 kind 不匹配导致无法累加) |
事件时间模式下各类丢弃(缺失时间戳、过于未来、迟到)统一计入通用的component_discarded_events_total计数器,由 AggregateEventDropped 触发。
使用建议与限制
- 聚合以降低时间粒度为代价换取数据量下降,
interval_ms越大压缩比越高、延迟与粒度损失越大; - 该组件有状态:进程内保留各序列当前值与事件时间桶,配置校验会拦截非法的毫秒参数;
- 若下游以指标时间戳为键写入(覆盖式更新),建议启用
event_time块并配合allowed_lateness_ms容忍上游乱序,用max_future_ms拦截时钟异常; - 更多用例可参考 tests.rs 中的测试场景与 组件元数据 中的输入输出示例。
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考