从 EWMA 指数移动平均库到 vmctl 进度条:原理、实现与在 VictoriaMetrics 中的应用
【免费下载链接】VictoriaMetricsVictoriaMetrics: fast, cost-effective monitoring solution and time series database项目地址: https://gitcode.com/GitHub_Trending/vi/VictoriaMetrics
指数移动平均(Exponentially Weighted Moving Average,EWMA)是一种以极低计算与内存成本持续跟踪序列“近期中心趋势”的经典算法。VictoriaMetrics 仓库通过 go.mod 以间接依赖方式引入了第三方库github.com/VividCortex/ewmav1.2.0(见 vendor/github.com/VividCortex/ewma 目录下的 vendored 副本),本文以该库的 README 与 ewma.go 源码为主体,讲透 EWMA 的算法原理、alpha(衰减因子)的选取方法、两种实现的取舍,并追踪它在 vmctl 迁移工具进度条速度估算中的真实调用链。
一、什么是指数移动平均
按 README 的定义:EWMA 是一种随着数值逐个到达而连续计算的平均方式。每个新样本加入后,其权重随时间指数衰减,使平均值偏向更近的历史数据。它的核心优势有三点:
- 计算成本低——每来一个样本只需一次乘加;
- 内存成本低——只需保存当前平均值一个状态量;
- 语义清晰——它表达的是序列“近期的中心趋势”,而非全历史平均。
算法本身可以用如下伪代码完整描述(继承自原文档):
- 将序列中的下一个数乘以 alpha;
- 将当前的平均值乘以 (1 - alpha);
- 将两步结果相加,存为新的当前平均值;
- 对序列中的每个数重复以上过程。
其中 alpha 是衰减因子,必须位于 (0, 1) 区间。alpha 越大,平均值对近期历史的偏向越强;实际中通常取一个较小的数(原文档举例 0.04)。
初始化这一“特例”
不同实现对“初始值怎么来”的处理并不一致。常见两种策略:
- 直接取第一个样本作为初始平均值(SimpleEWMA 的路线);
- 先对前 10 个左右样本做算术平均,再开始增量更新(VariableEWMA 的路线,来源是 Steven Nahmias 的《Production and Operations Analysis》的建议)。
两种各有利弊:前者零预热、更省内存,但第一个样本会被高估权重;后者让平均值“冷启动”更稳,代价是需要额外状态和一段预热期。
从直观上看,EWMA 是原序列的平滑(低通)平均:README 中用 α=0.5 的五点序列示意,每加入一个新值,旧样本对当前均值的贡献就减半,最终得到的是一条明显平滑于原始波动的曲线。
二、如何选取 alpha:alpha = 2/(N+1)
这是 README 中最有实操价值的一节。推理链条如下:
考虑一个固定窗口滑动平均,它对前 N 个样本求平均,那么这批样本的“平均年龄”恰好是 N/2;
若希望构造一个 EWMA,使其样本的平均年龄与上述 N 窗口平均相同,则所需的 alpha 满足公式:
alpha = 2 / (N + 1)
(该公式的证明见 Steven Nahmias 的《Production and Operations Analysis》。)
一个典型换算:若时间序列每秒采样一次,你希望得到“过去一分钟”的等效移动平均,则 N = 60,应取 alpha ≈0.032786885。
vendored 源码中的默认值
这一点可以直接对照 ewma.go 顶部的常量来验证:
const ( // By default, we average over a one-minute period, which means the average // age of the metrics in the period is 30 seconds. AVG_METRIC_AGE float64 = 30.0 // The formula for computing the decay factor from the average age comes // from "Production and Operations Analysis" by Steven Nahmias. DECAY float64 = 2 / (float64(AVG_METRIC_AGE) + 1) // ... WARMUP_SAMPLES uint8 = 10 )可以看到库作者把“平均样本年龄”定为30 秒(即一分钟窗口的 N/2),于是默认衰减因子DECAY = 2/31 ≈ 0.0645——正好落在一分钟窗口平均年龄的公式上。WARMUP_SAMPLES = 10则是 VariableEWMA 的预热样本数,对应 README 中“先对前 10 个样本做算术平均”的描述。
三、两种实现:SimpleEWMA 与 VariableEWMA
README 指出,库提供两种实现,二者都实现同一个MovingAverage接口,构造函数也返回该接口类型:
// ewma.go L28-L32 type MovingAverage interface { Add(float64) Value() float64 Set(float64) }还有一个重要的前提约束:所有实现都假设相邻两次 Add 之间的隐含时间间隔恒为 1.0,即把“时间流逝”等同于“样本到达”。如果你需要在采样间隔不规则时按真实时间做衰减,README 明确说明该包目前不满足这种需求。
SimpleEWMA:零预热、单字段、把 0 当未初始化
从源码看(ewma.go L59-L72),整个结构体只有一个字段:
type SimpleEWMA struct { value float64 } func (e *SimpleEWMA) Add(value float64) { if e.value == 0 { // this is a proxy for "uninitialized" e.value = value } else { e.value = (value * DECAY) + (e.value * (1 - DECAY))) } }关键行为:
- 无预热期,第一个样本直接成为平均值(
e.value == 0被当作“未初始化”的代理条件); - 常量衰减DECAY = 2/31;
- 零值陷阱:README 专门警告,若被平均的量在运行中真实地趋于 0,那么随后任何一个非零值都会造成“尖峰跳变”而不是小幅变化——因为 0 被解释成未初始化。README 同时补充:这个“衰减回零”的过程非常缓慢,值通常会稳定在一个接近 0 而非恰好为 0 的数值上,因此一般不会被误判为未初始化。
VariableEWMA:自定义年龄 + 预热期,约两倍内存
对比源码(ewma.go L86-L126):
type VariableEWMA struct { decay float64 // 衰减因子 2/(age+1) value float64 // 当前平均值 count uint8 // 已加入的样本数(预热阶段计数) } func (e *VariableEWMA) Add(value float64) { switch { case e.count < WARMUP_SAMPLES: // 前 10 个样本:累加 e.count++ e.value += value case e.count == WARMUP_SAMPLES: // 第 11 个样本:先取算术平均作为种子,再开始指数更新 e.count++ e.value = e.value / float64(WARMUP_SAMPLES) e.value = (value * e.decay) + (e.value * (1 - e.decay)) default: // 之后:标准 EWMA 递推 e.value = (value * e.decay) + (e.value * (1 - e.decay)) } }与 SimpleEWMA 的差异可以归纳为三点,与 README 描述一一对应:
- 支持自定义 age:需要持久保存
decay字段,因此内存更大(README 估计约为 SimpleEWMA 的两倍多); - 有预热期:前 10 个样本只做算术累加,
Value()在预热完成前始终返回 0.0(见 L112-L118 的if e.count <= WARMUP_SAMPLES { return 0.0 });这比 SimpleEWMA 的“首样本即平均值”冷启动更稳; - Set() 会强制越过预热期:
Set()在赋值后把count提升到WARMUP_SAMPLES + 1,意味着外部注入的值立即被视为可信。
构造函数:一行代码决定用哪个实现
NewMovingAverage 的分支逻辑非常值得注意:
func NewMovingAverage(age ...float64) MovingAverage { if len(age) == 0 || age[0] == AVG_METRIC_AGE { return new(SimpleEWMA) } return &VariableEWMA{ decay: 2 / (age[0] + 1), } }- 不传参数,或传入的 age 恰好等于默认的 30,直接返回
SimpleEWMA(最小内存); - 传其他 age 则返回
VariableEWMA,衰减因子按公式2/(age+1)计算。
age 的语义是“当时间趋于无穷时,样本的平均年龄”(README 原话)。
四、API 使用示例
README 给出了最小用法,这里按 vendored 源码的语义整理为可运行形式(具体输出数值取决于输入序列,此处只展示 API 形态):
package main import "github.com/VividCortex/ewma" func main() { samples := [100]float64{ 4599, 5711, 4746, 4621, 5037, 4218, 4925, 4281, 5207, 5203, 5594, 5149, } e := ewma.NewMovingAverage() // 无参 => SimpleEWMA,DECAY = 2/31 a := ewma.NewMovingAverage(5) // => VariableEWMA,decay = 2/(5+1) = 1/3 for _, f := range samples { e.Add(f) a.Add(f) } // 注意:a 只 Add 了 12 个样本(10 个预热 + 2 次更新), // 而 Value() 在 count <= 10 时返回 0.0,因此此处已可读到有效值。 e.Value() // 围绕样本量级(数千)的平滑均值 a.Value() }两点使用须知(均来自 README 与源码,而非推测):
VariableEWMA.Value()在加满 10 个样本之前恒为 0,依赖其输出做判断(例如“速率为 0 则显示 ?”)的调用方必须意识到这一点;- 对真实可能为 0 的量,优先选
VariableEWMA或确保值不会精确归零,以规避 SimpleEWMA 的零值跳变问题。
五、它在 VictoriaMetrics 中用在哪:vmctl 进度条的“p/s 速度”
在本仓库中,ewma并不是核心时序引擎的直接依赖——go.mod 将其标注为// indirect。从源码结构看,它的实际消费方是 vendored 的进度条库github.com/cheggaaa/pb/v3,用于估算并平滑显示处理速度。调用链如下:
vmctl 各迁移入口 (app/vmctl/prometheus.go 等) └─> barpool (app/vmctl/barpool/pool.go) // 全局进度条池 └─> cheggaaa/pb/v3 (vendor) // 模板渲染 {{speed .}} └─> VividCortex/ewma // 速度的指数移动平均速度估算的实现细节
在 vendor/github.com/cheggaaa/pb/v3/speed.go 中:
speed结构体持有一个ewma.MovingAverage(默认即NewMovingAverage(),也就是 30 秒平均年龄的 SimpleEWMA,见 L20-L23);- 每次进度条状态更新时,若距上次采样不足
speedAddLimit = time.Second / 2(0.5 秒,见 L11),则直接返回旧均值,避免高频抖动; - 满足间隔后,计算该区间的瞬时速度
diff / dur.Seconds(),Add进 EWMA,再取Value()作为“当前速度”(L34-L44); - 进度条模板中的
{{speed .}}由ElementSpeed渲染成类似1.2k p/s的输出(L72-L83),速率为 0 时显示? p/s——这一分支正是对 EWMA 冷启动/零值语义的兜底处理。
在 vmctl 中的落点
vmctl 是所有数据迁移入口(Prometheus、InfluxDB、OpenTSDB、remote read、VM 原生格式等),大量使用进度条反馈迁移进度。例如 app/vmctl/prometheus.go L133-L137:
bar := barpool.AddWithTemplate(fmt.Sprintf(barTpl, "Processing blocks"), len(blocks)) if err := barpool.Start(); err != nil { ... } defer barpool.Stop()app/vmctl/barpool/pool.go 进一步封装了全局进度条池(pb.NewPool()),支持终端/非终端两种渲染模式(getTemplate会判断stdout是否为终端来决定换行策略),并提供progressBarNoOp以便通过Disable(true)将进度条完全置空。也就是说,当你运行 vmctl 迁移时看到的“Processing blocks … p/s”速度数字,其平滑背后就是本库的 EWMA。
六、选型与局限小结
结合 README 与 vendored 源码,可以把这套库的适用边界概括为:
| 维度 | SimpleEWMA | VariableEWMA |
|---|---|---|
| 状态字段 | 1 个float64 | decay+value+count(约 2 倍多内存) |
| 衰减因子 | 固定 2/31(平均年龄 30s) | 自定义2/(age+1) |
| 冷启动 | 首样本即平均值 | 前 10 样本算术平均,之前Value()返回 0 |
| 零值语义 | 0 视为未初始化,真实归零后非零值会跳变 | 无此问题 |
| 时间基准 | 隐含样本间隔恒为 1.0 | 同左(两者都不支持真实时间衰减) |
几条实操建议(均可从源码确认):
- 高基数、要求极致省内存的场景用
NewMovingAverage()(SimpleEWMA),但要确认被平均的量不会真实归零; - 需要“平均窗口 = N 个样本”的语义、或量本身可能为 0 的场景,用
NewMovingAverage(N)(VariableEWMA),并接受其预热期内的 0 输出; - 依赖
Value()做判空/判断时,务必考虑 VariableEWMA 预热期返回 0.0 的行为; - 若采样间隔不规则且需要按墙钟时间衰减,此库不适用(README 明确声明),需要在库外自行实现。
参考文件索引:README、ewma.go、pb/v3 速度估算、vmctl 进度条池、vmctl Prometheus 迁移。
【免费下载链接】VictoriaMetricsVictoriaMetrics: fast, cost-effective monitoring solution and time series database项目地址: https://gitcode.com/GitHub_Trending/vi/VictoriaMetrics
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考