学习目标
学完本章,你应该能够:
- 讲清"预测式扩容"相比原生 HPA 解决了什么问题——为什么"事后补救"会在大促式脉冲流量下雪崩。
- 把一条完整的预测扩容链路拆成 6 个环节(采集 → 特征 → 训练 → 推理 → 决策 → 执行),说清每环职责与典型坑。
- 解释 Prophet + XGBoost 集成的思路:Prophet 抓周期/节假日基线,XGBoost 学残差,最终预测 = 基线 + 残差。
- 看懂扩容决策器的核心公式与工程护栏(预热时间、冷却窗口、缩容幅度、阈值留余量),并知道这些参数为什么是这么设的。
- 设计"预测器 + HPA"的协同分工,以及用 MAPE 漂移检测触发重训的 MLOps 闭环。
前置知识:
- Kubernetes 基础:Deployment、HPA、client-go 基本认知。
- Python 基础(Pandas、scikit-learn 概念),能读模型训练脚本。
- Go 基础,能读决策器 / 控制器代码。
- 上一篇笔记的 client-go 执行器(本文扩容执行器与其共用代码)。
本章你会动手做的事:
- 用文末的
fetch_metrics.py从 Prometheus 拉 30 天 QPS,跑通特征工程生成qps_features.csv。 - 本地起 FastAPI 推理服务,调
/predict拿到未来 15 分钟预测 QPS。 - 跑一遍扩容决策器
Decide,故意把预测 QPS 调到阈值之上,观察它算出"提前扩容"的目标副本数。
类比:原生 HPA 像"着火了才叫消防队"——CPU 已经飙高了才扩容,等新车(Pod)开到现场火都烧半天了。预测式扩容像"看天气预报带伞"——模型说下午有暴雨(流量高峰),你上午就把伞(副本)准备好。本文就是把"看预报带伞"这套系统从数据采集到 K8s 执行完整落地。
实战背景
说实话我一开始觉得 HPA 够用了——CPU 高了扩容,CPU 低了缩容,多简单。直到有一次大促,流量从 500 QPS 飙到 5000 QPS 只用了 30 秒,HPA 从检测到 CPU 飙高到新 Pod Ready 总共花了 90 秒——这 90 秒里 P99 延迟飙到 8 秒,用户疯狂刷新,雪崩了。那天晚上我加班到凌晨三点,第二天就开始研究预测性扩容。
Kubernetes 原生的 HPA 基于当前指标(如 CPU、内存、自定义指标)做扩缩容,本质上是"事后补救"——指标已经飙高了才触发扩容,等新 Pod 起来流量早就冲过来了。对于流量具有明显周期性的业务(如电商促销、早晚高峰、直播带货),如果能提前几分钟甚至几小时预测流量峰值,在高峰到来之前就把副本扩好,就能从"被动挨打"变成"提前布防"。
这篇笔记记录的是我从零搭建一套预测性扩容系统的过程。模型不复杂(Prophet + XGBoost 集成),但工程链路很长——从 Prometheus 拉数据、特征工程、模型训练、推理服务、到扩容决策器、再到 K8s 执行,每一环都有坑。
整体架构
类比:整套系统分"白天想对策"和"实时执行"两条线。离线训练像军师在后方研究历史战报、练出一套预测兵法(模型);在线推理 + 扩容像前线指挥,每分钟看一眼当前敌情(实时指标),用兵法推演未来 30 分钟,提前调兵(扩副本)。HPA 则是督战官,负责军师没料到的突发状况兜底。
下面这张图把"离线训练"和"在线推理 + 扩容"两条链路一次性画清楚:
flowchart LR
subgraph 离线训练
P1[Prometheus 历史 QPS] --> P2[特征工程]
P2 --> P3[模型训练
Prophet + XGBoost]
P3 --> P4[模型注册/下载]
end
subgraph 在线推理与扩容
R1[当前指标 Prometheus] --> R2[推理服务 FastAPI]
P4 --> R2
R2 --> R3[预测未来 30 分钟 QPS]
R3 --> R4[扩容决策器 Go]
R4 --> R5[K8s API 改副本]
R6[HPA 兜底
基于实时 CPU] -. 精细微调 .-> R5
end┌─────────────────────────────────────────────────────────────────────┐
│ 离线训练流程 │
│ │
│ Prometheus ──▶ 历史指标采集 ──▶ 特征工程 ──▶ 模型训练 ──▶ 模型注册 │
│ (30天 QPS) (Python 脚本) (Pandas) (Prophet/XGB) (MLflow) │
└─────────────────────────────────────────────────────────────────────┘
│
模型文件下载
│
▼
┌─────────────────────────────────────────────────────────────────────┐
│ 在线推理 + 扩容流程 │
│ │
│ 当前指标 ──▶ 推理服务 ──▶ 预测未来 QPS ──▶ 扩容决策器 ──▶ K8s API │
│ (Prometheus) (FastAPI) (未来30分钟) (Go 程序) (client-go)│
│ │
│ ┌──────────────┐ │
│ │ HPA 兜底 │ ← 基于实时 CPU 做精细调整 │
│ └──────────────┘ │
└─────────────────────────────────────────────────────────────────────┘
主要组件的职责:
- 指标采集模块:从 Prometheus 拉取历史 QPS/CPU/内存数据,写成 CSV 或直接存入时序数据库。
- 特征工程模块:处理时间序列,生成训练样本。这是决定模型效果的关键步骤。
- 模型训练模块:Prophet 做基线预测,XGBoost 做残差修正,两个模型集成。
- 推理服务:FastAPI 部署,每分钟返回未来 30 分钟的预测 QPS。
- 扩容决策器:Go 程序,每分钟调推理服务,根据预测结果算目标副本数。
- 执行器:调 K8s API 修改 Deployment 副本数,跟上一篇笔记的 client-go 执行器共用代码。
数据准备与特征工程
从 Prometheus 拉取历史数据
先写个脚本把 Prometheus 里的 QPS 数据拉出来。这个脚本看着简单,但第一次跑的时候我踩了个大坑——Prometheus 默认只存 15 天数据,我想要 30 天的历史数据发现拉不到,得改 retention 配置。
# scripts/fetch_metrics.py
# 从 Prometheus 拉取历史 QPS 数据,保存为 CSV
# 运行:python fetch_metrics.py --prom-url http://prometheus:9090 --days 30
import requests
import pandas as pd
import argparse
from datetime import datetime, timedelta
import time
def fetch_qps_from_prometheus(prom_url, query, start_time, end_time, step='60s'):
"""从 Prometheus 拉取 range query 数据
Args:
prom_url: Prometheus 地址
query: PromQL 查询语句
start_time: 开始时间 (datetime)
end_time: 结束时间 (datetime)
step: 采样间隔,默认 60 秒
Returns:
DataFrame,包含 timestamp 和 value 两列
"""
# Prometheus range query API
# 一次最多拉 11000 个数据点,超过需要分段
url = f"{prom_url}/api/v1/query_range"
params = {
'query': query,
'start': start_time.timestamp(),
'end': end_time.timestamp(),
'step': step,
}
resp = requests.get(url, params=params, timeout=60)
resp.raise_for_status()
data = resp.json()
if data['status'] != 'success':
raise ValueError(f"Prometheus query failed: {data}")
result = data['data']['result']
if len(result) == 0:
raise ValueError("No data returned from Prometheus")
# 取第一个时间序列(如果有多条,后面再聚合)
values = result[0]['values']
df = pd.DataFrame(values, columns=['timestamp', 'value'])
# Prometheus 返回的是 Unix 时间戳(字符串),需要转换
df['timestamp'] = pd.to_datetime(df['timestamp'].astype(int), unit='s')
# value 是字符串,转成 float
df['value'] = df['value'].astype(float)
df = df.set_index('timestamp')
return df
def fetch_history(prom_url, days=30):
"""拉取 N 天的 QPS 数据,分段拉避免单次请求过大"""
end_time = datetime.now()
start_time = end_time - timedelta(days=days)
# 用 sum(rate()) 算总 QPS
# 注意:rate() 的窗口 [1m] 要跟采样间隔匹配
# 如果服务有多个实例,用 sum 聚合
query = 'sum(rate(http_requests_total{job="api-gateway"}[1m]))'
all_dfs = []
# 每次拉 1 天,避免单次请求超时
current = start_time
while current < end_time:
chunk_end = min(current + timedelta(days=1), end_time)
print(f"fetching {current} to {chunk_end}...")
try:
df = fetch_qps_from_prometheus(prom_url, query, current, chunk_end)
all_dfs.append(df)
except Exception as e:
print(f" warning: failed to fetch this chunk: {e}")
# 跳过失败的段,不要让一个 chunk 失败导致整个拉取失败
current = chunk_end
time.sleep(1) # 别把 Prometheus 打爆
if not all_dfs:
raise RuntimeError("No data fetched")
result = pd.concat(all_dfs)
# 去重(分段拉取可能有重叠)
result = result[~result.index.duplicated(keep='first')]
result = result.sort_index()
return result
if __name__ == '__main__':
parser = argparse.ArgumentParser()
parser.add_argument('--prom-url', default='http://prometheus:9090')
parser.add_argument('--days', type=int, default=30)
parser.add_argument('--output', default='qps_history.csv')
args = parser.parse_args()
df = fetch_history(args.prom_url, args.days)
df.to_csv(args.output)
print(f"saved {len(df)} rows to {args.output}")
print(f"date range: {df.index[0]} to {df.index[-1]}")
print(f"qps range: {df['value'].min():.1f} - {df['value'].max():.1f}")
踩坑提示:
rate(metric[1m])的窗口要跟你的采样间隔匹配。如果采样间隔 60 秒但 rate 窗口设 5m,数据会有冗余平滑。- Prometheus range query 单次返回上限是 11000 点。30 天 × 1440 分钟/天 = 43200 点,必须分段拉。
http_requests_total是 counter 类型,只增不减。如果服务重启了 counter 会 reset,导致 rate() 出现负值或尖峰。Prometheus 内部会处理这个问题,但如果你用 raw counter 自己算差值就会踩坑。- 如果你的服务有多个 Pod,
sum(rate(...))一定要在最外层 sum,不要先 sum 再 rate——counter sum 之后再 rate 会出错。
特征工程
时序预测的核心是构造高质量的训练数据。我试过很多特征组合,最后留下来的是这几个:
| 特征类型 | 示例 | 为什么有用 |
|---|---|---|
| 时间特征 | 小时、星期、是否节假日 | 流量有明显的日周期和周周期 |
| 滞后特征 | 过去 1h、过去 24h 的 QPS | 最近的历史值是最强的预测信号 |
| 滑动统计 | 过去 7 天同一时刻均值、标准差 | 捕捉周期性模式,比单点滞后更稳定 |
| 外部特征 | 营销活动标记、版本发布标记 | 促销日流量可能翻 10 倍,不标记模型会懵 |
# scripts/feature_engineering.py
# 特征工程:从原始 QPS 数据生成训练特征
import pandas as pd
import numpy as np
from datetime import datetime
def build_features(df):
"""从原始 QPS DataFrame 生成特征矩阵
Args:
df: DataFrame,index 是 datetime,有一列 'value' 是 QPS
Returns:
DataFrame,包含原始 QPS + 所有特征
"""
df = df.copy()
df = df.rename(columns={'value': 'qps'})
# === 时间特征 ===
# hour: 0-23,捕捉日内周期
df['hour'] = df.index.hour
# dayofweek: 0=周一, 6=周日
df['dayofweek'] = df.index.dayofweek
# is_weekend: 周末流量模式跟工作日不一样
df['is_weekend'] = (df.index.dayofweek >= 5).astype(int)
# 节假日特征:需要外部数据源
# 我用的是中国法定节假日,自己维护一个日期列表
# 如果有营销活动也在这里加
holidays = [
'2025-01-01', '2025-02-10', '2025-02-11', '2025-02-12', # 春节
'2025-04-04', '2025-04-05', '2025-04-06', # 清明
'2025-05-01', '2025-05-02', '2025-05-03', # 劳动节
'2025-06-14', '2025-06-15', '2025-06-16', # 端午
'2025-09-27', '2025-09-28', '2025-09-29', # 中秋
'2025-10-01', '2025-10-02', '2025-10-03', # 国庆
# 电商大促
'2025-06-18', '2025-11-11', '2025-12-12',
]
holiday_dates = pd.to_datetime(holidays).date
df['is_holiday'] = df.index.date.isin(holiday_dates).astype(int)
# === 滞后特征 ===
# 1 小时前的 QPS:短期趋势信号
df['qps_lag_1h'] = df['qps'].shift(1)
# 24 小时前的 QPS:昨天同一时刻的 QPS,最强的周期信号
df['qps_lag_24h'] = df['qps'].shift(24)
# 7 天前同一时刻的 QPS:上周同一天的 QPS
df['qps_lag_7d'] = df['qps'].shift(7 * 24)
# === 滑动统计特征 ===
# 过去 7 天同一时刻的均值和标准差
# 用 shift(24) 先拿到"昨天同一时刻",再 rolling 7 天
# 这样每个窗口包含的是过去 7 天的同一时刻,而不是过去 7 天的所有时刻
df['qps_same_hour_mean_7d'] = df['qps'].shift(24).rolling(window=7).mean()
df['qps_same_hour_std_7d'] = df['qps'].shift(24).rolling(window=7).std()
# 过去 1 小时的均值和变化率
df['qps_roll_mean_1h'] = df['qps'].rolling(window=60).mean() # 60 个 1 分钟点
df['qps_roll_std_1h'] = df['qps'].rolling(window=60).std()
# QPS 变化率(一阶差分):捕捉趋势
df['qps_diff_1'] = df['qps'].diff(1)
df['qps_diff_5'] = df['qps'].diff(5)
# === 异常值处理 ===
# QPS 不应该有负值或极端大值
# 用 3σ 法则裁剪,但别直接删,用 clip 限幅
qps_mean = df['qps'].mean()
qps_std = df['qps'].std()
upper_bound = qps_mean + 5 * qps_std # 5σ 而不是 3σ,避免裁掉正常峰值
df['qps'] = df['qps'].clip(lower=0, upper=upper_bound)
# 删除因为有 NaN 的行(前 7*24=168 行因为 lag 特效必然是 NaN)
df = df.dropna()
return df
if __name__ == '__main__':
df = pd.read_csv('qps_history.csv', parse_dates=['timestamp'], index_col='timestamp')
features = build_features(df)
features.to_csv('qps_features.csv')
print(f"generated {len(features.columns)} features, {len(features)} rows")
print(f"features: {list(features.columns)}")
踩坑提示:
qps_lag_24h这个特征权重通常最高——昨天同一时刻的 QPS 是最强的预测信号。但如果你的服务是刚上线的,没有 24 小时前的数据,这个特征就是 NaN,模型直接不可用。新服务至少要跑 7 天再开始训练。- 节假日列表要手动维护。我一开始忘了加双十一标记,模型在 11 月 11 日的预测偏差了 8 倍——它按正常周四的流量预测的。
rolling(window=7)在shift(24)之后做,这个顺序很重要。如果反过来先 rolling 再 shift,拿到的就是"过去 7 天全部时刻的均值",周期信号就被抹平了。- 异常值裁剪用
clip而不是删除行。删行会导致时间序列断裂,lag 特征会错位。
模型选择与训练
模型对比
根据数据规模和业务特点,可选择不同的时序预测模型。我把几个用过的模型横向对比一下:
| 模型 | 适用场景 | 我的评价 |
|---|---|---|
| ARIMA / SARIMA | 数据量小、趋势和季节性明显 | 参数调起来麻烦(p/d/q),而且非线性的突发流量完全预测不了 |
| Prophet | 需要解释性强、包含节假日效应 | 开箱即用,节假日效应内置,但单变量预测精度有限 |
| LSTM / GRU | 数据量大、非线性关系复杂 | 效果好但训练慢,还要调超参,生产部署也麻烦 |
| XGBoost / LightGBM | 特征工程丰富、需要快速训练 | 我的最终选择,训练快、特征重要性可解释、精度够用 |
| Transformers(如 PatchTST) | 长序列、高精度需求 | 太重了,不值得,QPS 预测用不着这么复杂的模型 |
我最终用的是 Prophet + XGBoost 集成:Prophet 做基线预测(擅长捕捉日周期和节假日效应),XGBoost 做残差修正(用工程化特征补充 Prophet 捕捉不到的非线性模式)。
Prophet 基线模型
# scripts/train_prophet.py
# 训练 Prophet 模型,做基线预测
import pandas as pd
from prophet import Prophet
import joblib
import logging
# Prophet 的日志太吵了,设为 WARNING
logging.getLogger('prophet').setLevel(logging.WARNING)
logging.getLogger('cmdstanpy').setLevel(logging.WARNING)
def train_prophet(df):
"""训练 Prophet 模型
Prophet 要求输入 DataFrame 有两列:
- ds: 日期时间
- y: 目标值
Args:
df: 特征工程后的 DataFrame,包含 qps 列
Returns:
训练好的 Prophet 模型
"""
# Prophet 只需要 ds 和 y 两列,其他特征它自己会算季节性
train_df = pd.DataFrame({
'ds': df.index,
'y': df['qps'].values,
})
# 初始化 Prophet
# yearly_seasonality=False: 一年的数据不够学年度季节性
# weekly_seasonality=True: 周周期很重要,默认开
# daily_seasonality=True: 日周期是核心,默认开
# changepoint_prior_scale: 控制趋势变化的灵活度,默认 0.05
# 调大到 0.1 让模型更敏感地捕捉趋势变化(但别太大,会过拟合)
model = Prophet(
yearly_seasonality=False,
weekly_seasonality=True,
daily_seasonality=True,
changepoint_prior_scale=0.1,
changepoint_range=0.9, # 用前 90% 的数据来检测趋势变化点
)
# 添加节假日效应
# Prophet 内置了部分国家的节假日,但中国节假日要自己加
holidays = pd.DataFrame({
'holiday': 'cn_holiday',
'ds': pd.to_datetime([
'2025-01-01', '2025-02-10', '2025-02-11', '2025-02-12',
'2025-04-04', '2025-05-01', '2025-06-18',
'2025-09-27', '2025-10-01', '2025-10-02', '2025-10-03',
'2025-11-11', '2025-12-12',
]),
'lower_window': -1, # 节假日前 1 天也开始有影响
'upper_window': 1, # 节假日后 1 天还有余温
})
# Prophet 内置了部分国家的节假日,但中国节假日不一定有
# 尝试加载内置 CN 节假日,失败就用上面手动定义的 holidays
try:
model = model.add_country_holidays(country_name='CN')
except Exception as e:
print(f"no built-in CN holidays, using manual list: {e}")
# 直接传 holidays 作为补充(跟内置节假日不冲突,Prophet 会合并)
model.holidays = holidays
print("training Prophet model...")
model.fit(train_df)
return model
if __name__ == '__main__':
df = pd.read_csv('qps_features.csv', parse_dates=['timestamp'], index_col='timestamp')
model = train_prophet(df)
# 保存模型
joblib.dump(model, 'prophet_model.pkl')
print("model saved to prophet_model.pkl")
# 快速验证:预测未来 30 分钟
future = model.make_future_dataframe(periods=30, freq='min')
forecast = model.predict(future)
print("\n=== forecast for next 30 min ===")
print(forecast[['ds', 'yhat', 'yhat_lower', 'yhat_upper']].tail(30).to_string(index=False))
XGBoost 残差修正模型
Prophet 的基线预测有系统性偏差——它对突发流量的反应比较慢。用 XGBoost 学习这个残差:
# scripts/train_xgb.py
# 训练 XGBoost 模型,修正 Prophet 的预测残差
import pandas as pd
import numpy as np
import joblib
from prophet import Prophet
from xgboost import XGBRegressor
from sklearn.model_selection import TimeSeriesSplit
from sklearn.metrics import mean_absolute_error, mean_absolute_percentage_error
import logging
logging.getLogger('prophet').setLevel(logging.WARNING)
def train_residual_model(df, prophet_model):
"""用 XGBoost 学习 Prophet 预测的残差
策略:先用 Prophet 预测训练集,算出残差(实际值 - 预测值),
然后用 XGBoost 学习残差与特征之间的关系。
最终预测 = Prophet 预测 + XGBoost 预测的残差
"""
# 1. 用 Prophet 预测训练集(in-sample 预测)
prophet_input = pd.DataFrame({
'ds': df.index,
'y': df['qps'].values,
})
prophet_pred = prophet_model.predict(prophet_input)
# 2. 计算残差
residual = df['qps'].values - prophet_pred['yhat'].values
# 3. 准备 XGBoost 的特征和标签
# 用工程化特征(hour, dayofweek, lag 特征等)来预测残差
feature_cols = [
'hour', 'dayofweek', 'is_weekend', 'is_holiday',
'qps_lag_1h', 'qps_lag_24h', 'qps_lag_7d',
'qps_same_hour_mean_7d', 'qps_same_hour_std_7d',
'qps_roll_mean_1h', 'qps_roll_std_1h',
'qps_diff_1', 'qps_diff_5',
]
# 加上 Prophet 的预测值作为特征(让 XGB 知道基线是多少)
df_features = df[feature_cols].copy()
df_features['prophet_yhat'] = prophet_pred['yhat'].values
X = df_features.values
y = residual # 标签是残差,不是原始 QPS
# 4. 时间序列交叉验证
# 不能用随机 K-Fold,必须按时间切分
tscv = TimeSeriesSplit(n_splits=5)
best_mae = float('inf')
for fold, (train_idx, val_idx) in enumerate(tscv.split(X)):
X_train, X_val = X[train_idx], X[val_idx]
y_train, y_val = y[train_idx], y[val_idx]
model = XGBRegressor(
n_estimators=300, # 树的数量,300 足够
max_depth=6, # 深度 6 防过拟合
learning_rate=0.05, # 小学习率 + 多树 = 更稳定
subsample=0.8, # 行采样
colsample_bytree=0.8, # 列采样
reg_alpha=0.1, # L1 正则化
reg_lambda=1.0, # L2 正则化
random_state=42,
n_jobs=-1,
)
model.fit(
X_train, y_train,
eval_set=[(X_val, y_val)],
verbose=False,
)
# 验证集预测
val_pred = model.predict(X_val)
mae = mean_absolute_error(y_val, val_pred)
# 最终预测 = Prophet 基线 + XGB 残差
final_pred = prophet_pred['yhat'].values[val_idx] + val_pred
final_mae = mean_absolute_error(df['qps'].values[val_idx], final_pred)
final_mape = mean_absolute_percentage_error(
df['qps'].values[val_idx] + 1, # +1 避免除零
final_pred + 1,
)
print(f"fold {fold}: residual MAE={mae:.1f}, final MAE={final_mae:.1f}, MAPE={final_mape:.1%}")
if final_mae < best_mae:
best_mae = final_mae
# 用全部数据重新训练
final_model = XGBRegressor(
n_estimators=300, max_depth=6, learning_rate=0.05,
subsample=0.8, colsample_bytree=0.8,
reg_alpha=0.1, reg_lambda=1.0, random_state=42, n_jobs=-1,
)
final_model.fit(X, y, verbose=False)
return final_model, feature_cols, best_mae
if __name__ == '__main__':
df = pd.read_csv('qps_features.csv', parse_dates=['timestamp'], index_col='timestamp')
# 先加载 Prophet 模型
prophet_model = joblib.load('prophet_model.pkl')
# 训练 XGBoost 残差模型
xgb_model, feature_cols, best_mae = train_residual_model(df, prophet_model)
# 保存
joblib.dump({
'model': xgb_model,
'feature_cols': feature_cols,
}, 'xgb_residual_model.pkl')
print(f"\nmodel saved. best final MAE: {best_mae:.1f}")
踩坑提示:
TimeSeriesSplit不能换成普通的KFold。时序数据有严格的时间顺序,用随机切分会导致"用未来的数据预测过去",训练指标很好看但线上直接拉胯。- XGBoost 的
max_depth别超过 8。QPS 数据噪声大,树太深会过拟合——训练 MAE 很低但验证 MAE 飙高。 - 残差可能有负值,XGBoost 回归能处理,但要注意 Prophet 在低流量时段(凌晨 3 点)的预测经常偏高,残差偏负。如果你的模型在低流量时段预测为负 QPS,加一个
max(0, pred)后处理。 - 我试过用 LSTM 替代 XGBoost,效果没好多少(MAE 差 3%),但训练时间从 30 秒变成 2 小时,部署也复杂得多。QPS 预测这种单变量时序问题,XGBoost 足够了。
模型服务化部署
FastAPI 推理服务
训练好的模型要部署成服务,供扩容决策器调用。我用 FastAPI,因为简单、自带 Swagger、异步支持好:
# app/main.py
# FastAPI 推理服务:输入预测步数,返回未来 N 分钟的 QPS 预测
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
import joblib
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import logging
# 日志配置
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
app = FastAPI(title="QPS Predictor", version="1.0.0")
# 全局变量,启动时加载模型
prophet_model = None
xgb_model = None
xgb_feature_cols = None
@app.on_event("startup")
def load_models():
"""启动时加载模型,避免每次请求都加载"""
global prophet_model, xgb_model, xgb_feature_cols
prophet_model = joblib.load('models/prophet_model.pkl')
xgb_data = joblib.load('models/xgb_residual_model.pkl')
xgb_model = xgb_data['model']
xgb_feature_cols = xgb_data['feature_cols']
logger.info("models loaded successfully")
class PredictRequest(BaseModel):
steps: int = 30 # 预测未来多少分钟
class PredictPoint(BaseModel):
timestamp: str
qps: float
qps_lower: float
qps_upper: float
class PredictResponse(BaseModel):
predictions: list[PredictPoint]
model_version: str = "1.0"
@app.post("/predict", response_model=PredictResponse)
def predict(req: PredictRequest):
"""预测未来 N 分钟的 QPS
流程:
1. Prophet 做基线预测
2. 构造特征,XGBoost 预测残差
3. 最终预测 = Prophet 基线 + XGB 残差
4. 加上置信区间
"""
if req.steps <= 0 or req.steps > 120:
raise HTTPException(status_code=400, detail="steps must be between 1 and 120")
# 1. Prophet 基线预测
future = prophet_model.make_future_dataframe(periods=req.steps, freq='min')
prophet_forecast = prophet_model.predict(future)
# 取最后 N 分钟的预测
prophet_pred = prophet_forecast.tail(req.steps).copy()
# 2. 构造 XGBoost 特征
# 注意:预测时没有实际的 QPS 值,lag 特征要用 Prophet 的预测值来填充
# 这是一个简化处理——更严谨的做法是用递归预测(每次预测一步,用预测值作为下一步的 lag)
now = datetime.now()
features = []
for i, row in prophet_pred.iterrows():
ts = row['ds']
feat = {
'hour': ts.hour,
'dayofweek': ts.dayofweek,
'is_weekend': int(ts.dayofweek >= 5),
'is_holiday': 0, # 简化:预测时不知道未来是否节假日,实际要查日历
# lag 特征用 Prophet 预测值近似
# 严格来说应该用递归预测,但 Prophet 预测已经足够准了
'qps_lag_1h': row['yhat'],
'qps_lag_24h': row['yhat'],
'qps_lag_7d': row['yhat'],
'qps_same_hour_mean_7d': row['yhat'],
'qps_same_hour_std_7d': 0,
'qps_roll_mean_1h': row['yhat'],
'qps_roll_std_1h': 0,
'qps_diff_1': 0,
'qps_diff_5': 0,
'prophet_yhat': row['yhat'],
}
features.append(feat)
# 3. XGBoost 残差预测
feature_df = pd.DataFrame(features)
X = feature_df[xgb_feature_cols + ['prophet_yhat']].values
residual_pred = xgb_model.predict(X)
# 4. 最终预测 = 基线 + 残差
final_pred = prophet_pred['yhat'].values + residual_pred
# QPS 不能为负
final_pred = np.maximum(final_pred, 0)
# 置信区间:用 Prophet 的区间 + 残差的标准差
# 这是一个近似,更准确的做法是用分位数回归
residual_std = np.std(residual_pred)
lower = final_pred - 1.96 * residual_std
upper = final_pred + 1.96 * residual_std
# 构造响应
predictions = []
for i, (_, row) in enumerate(prophet_pred.iterrows()):
predictions.append(PredictPoint(
timestamp=row['ds'].isoformat(),
qps=round(final_pred[i], 1),
qps_lower=round(max(0, lower[i]), 1),
qps_upper=round(upper[i], 1),
))
logger.info(f"predicted {req.steps} steps, max_qps={max(p.qps for p in predictions):.1f}")
return PredictResponse(predictions=predictions)
@app.get("/health")
def health():
"""健康检查端点"""
if prophet_model is None or xgb_model is None:
raise HTTPException(status_code=503, detail="models not loaded")
return {"status": "ok"}
# 启动命令:uvicorn app.main:app --host 0.0.0.0 --port 8080
Dockerfile 和 K8s 部署:
# Dockerfile
FROM python:3.11-slim
WORKDIR /app
# 安装依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# 复制代码和模型
COPY app/ ./app/
COPY models/ ./models/
EXPOSE 8080
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8080"]
# k8s-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: qps-predictor
namespace: aiops
spec:
replicas: 2 # 两个副本做 HA
selector:
matchLabels:
app: qps-predictor
template:
metadata:
labels:
app: qps-predictor
spec:
containers:
- name: predictor
image: registry.example.com/qps-predictor:v1.0
ports:
- containerPort: 8080
resources:
requests:
cpu: 200m
memory: 512Mi
limits:
cpu: 500m
memory: 1Gi
# 健康检查
livenessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 10
periodSeconds: 30
readinessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 5
periodSeconds: 10
---
apiVersion: v1
kind: Service
metadata:
name: qps-predictor
namespace: aiops
spec:
selector:
app: qps-predictor
ports:
- port: 8080
targetPort: 8080
踩坑提示:
- 模型加载放到
@app.on_event("startup")里,别放在模块顶层。顶层加载会导致每次 import 都加载一次模型,测试时特别慢。 - Prophet 的
make_future_dataframe是基于模型训练时的最后时间戳往后推的,不是从当前时间推。所以如果你的模型是一天前训练的,预测出来的时间戳会差一天。要么每次重训,要么在请求时手动调整。 - 预测时
is_holiday那里我简化成了 0。实际要做的话,得提前把未来 N 天的节假日标记传进推理服务,或者让推理服务自己查日历 API。 uvicorn默认单进程。如果你的 QPS 比较高,用--workers 2起多进程。但注意每个 worker 都会加载一份模型,内存占用会翻倍。
扩容决策与执行
类比:扩容决策器像电梯的"载重控制器"。它不只看"现在有多少人",还看"接下来几站要上多少人"——所以你按了上行,电梯会提前在上一层就多停一会儿接人。本文决策器也是:不看所有预测点,只看"Pod 启动延迟之后"那段时间内的最大预测 QPS,提前把副本扩好。
扩容决策的核心判断流程如下:
flowchart TD
A[拿到预测窗口内各点 QPS] --> B[跳过 Pod 启动延迟内的点]
B --> C[取剩余点中的最大预测 QPS]
C --> D{在冷却期内?}
D -->|是| E[no_action 不扩缩]
D -->|否| F{最大QPS > 扩容阈值?}
F -->|是| G[scale_up 留 20% 余量]
F -->|否| H{最大QPS < 缩容阈值?}
H -->|是| I[scale_down 每次最多减 30%]
H -->|否| J[no_action 维持现状]决策算法
扩容决策不是简单的"预测 QPS / 单 Pod QPS = 目标副本数"。要考虑预热时间、冷却窗口、最小最大副本限制等多个因素:
// pkg/scaler/decision.go
package scaler
import (
"fmt"
"math"
"time"
)
// ScaleConfig 扩容配置
type ScaleConfig struct {
MinReplicas int32 // 最小副本数
MaxReplicas int32 // 最大副本数
QPSPerPod float64 // 每个 Pod 能承载的 QPS(压测得出)
ScaleUpThreshold float64 // 扩容阈值:预测 QPS 达到容量的多少百分比时扩容
ScaleDownThreshold float64 // 缩容阈值
CooldownSeconds int // 冷却时间:两次扩缩容之间最少间隔
PredictSteps int // 预测步数(分钟)
PodStartupSeconds int // Pod 启动时间(预热),用于提前扩容
}
// DefaultConfig 默认配置
func DefaultConfig() ScaleConfig {
return ScaleConfig{
MinReplicas: 3,
MaxReplicas: 50,
QPSPerPod: 200, // 压测得出:单 Pod 200 QPS 时 P99 < 200ms
ScaleUpThreshold: 0.7, // 达到 70% 容量就开始扩容,留余量
ScaleDownThreshold: 0.3, // 降到 30% 容量才缩容,避免抖动
CooldownSeconds: 120, // 2 分钟内不重复扩缩容
PredictSteps: 15, // 预测未来 15 分钟
PodStartupSeconds: 60, // Pod 启动 + 预热需要 60 秒
}
}
// ScaleDecision 扩缩容决策结果
type ScaleDecision struct {
CurrentReplicas int32 `json:"currentReplicas"`
DesiredReplicas int32 `json:"desiredReplicas"`
Reason string `json:"reason"`
PredictedMaxQPS float64 `json:"predictedMaxQPS"`
Action string `json:"action"` // scale_up / scale_down / no_action
}
// PredictedPoint 预测数据点
type PredictedPoint struct {
Timestamp string `json:"timestamp"`
QPS float64 `json:"qps"`
}
// Decide 根据预测结果和当前状态决定目标副本数
// 这是整个扩容系统的"大脑"
func Decide(
predictions []PredictedPoint,
currentReplicas int32,
lastScaleTime time.Time,
config ScaleConfig,
) ScaleDecision {
decision := ScaleDecision{
CurrentReplicas: currentReplicas,
}
if len(predictions) == 0 {
decision.Action = "no_action"
decision.Reason = "no prediction data"
decision.DesiredReplicas = currentReplicas
return decision
}
// 1. 找到预测窗口内的最大 QPS
// 不看所有预测点,只看"Pod 启动时间"之后的点
// 因为 Pod 启动需要时间,现在扩容要考虑预热
startupDelay := time.Duration(config.PodStartupSeconds) * time.Second
now := time.Now()
var maxPredictedQPS float64
for _, p := range predictions {
ts, err := time.Parse(time.RFC3339, p.Timestamp)
if err != nil {
continue
}
// 跳过 Pod 启动时间内的预测——这段时间来不及扩容
if ts.Sub(now) < startupDelay {
continue
}
if p.QPS > maxPredictedQPS {
maxPredictedQPS = p.QPS
}
}
if maxPredictedQPS == 0 {
maxPredictedQPS = predictions[0].QPS // fallback
}
decision.PredictedMaxQPS = maxPredictedQPS
// 2. 计算需要的副本数
// 公式:目标副本数 = ceil(预测 QPS / 单 Pod QPS)
// 但要考虑阈值:达到 70% 容量就扩容,而不是等到 100%
currentCapacity := float64(currentReplicas) * config.QPSPerPod
thresholdForScaleUp := currentCapacity * config.ScaleUpThreshold
// 3. 冷却期检查
timeSinceLastScale := time.Since(lastScaleTime)
if timeSinceLastScale < time.Duration(config.CooldownSeconds)*time.Second {
decision.Action = "no_action"
decision.Reason = fmt.Sprintf("cooldown: last scale %v ago, need %v",
timeSinceLastScale.Round(time.Second),
time.Duration(config.CooldownSeconds)*time.Second)
decision.DesiredReplicas = currentReplicas
return decision
}
// 4. 扩容判断
if maxPredictedQPS > thresholdForScaleUp {
// 需要扩容
desired := int32(math.Ceil(maxPredictedQPS / config.QPSPerPod))
// 留 20% 余量,避免刚扩完又得扩
desired = int32(math.Ceil(float64(desired) * 1.2))
desired = clampReplicas(desired, config.MinReplicas, config.MaxReplicas)
if desired > currentReplicas {
decision.Action = "scale_up"
decision.Reason = fmt.Sprintf("predicted max QPS %.0f > threshold %.0f, scaling from %d to %d",
maxPredictedQPS, thresholdForScaleUp, currentReplicas, desired)
decision.DesiredReplicas = desired
return decision
}
}
// 5. 缩容判断
// 缩容更保守:只有当预测 QPS 持续低于 30% 容量才缩
thresholdForScaleDown := currentCapacity * config.ScaleDownThreshold
if maxPredictedQPS < thresholdForScaleDown && currentReplicas > config.MinReplicas {
desired := int32(math.Ceil(maxPredictedQPS / config.QPSPerPod))
// 缩容别一次缩太多,每次最多缩 30%
maxReduction := int32(math.Ceil(float64(currentReplicas) * 0.3))
minDesired := currentReplicas - maxReduction
// 取 QPS 需求和最大减幅限制的较大值,避免缩过头
if desired < minDesired {
desired = minDesired
}
if desired < config.MinReplicas {
desired = config.MinReplicas
}
decision.Action = "scale_down"
decision.Reason = fmt.Sprintf("predicted max QPS %.0f < threshold %.0f, scaling down from %d to %d",
maxPredictedQPS, thresholdForScaleDown, currentReplicas, desired)
decision.DesiredReplicas = desired
return decision
}
// 6. 不需要调整
decision.Action = "no_action"
decision.Reason = fmt.Sprintf("predicted max QPS %.0f within [%.0f, %.0f]",
maxPredictedQPS, thresholdForScaleDown, thresholdForScaleUp)
decision.DesiredReplicas = currentReplicas
return decision
}
func clampReplicas(desired, min, max int32) int32 {
if desired < min {
return min
}
if desired > max {
return max
}
return desired
}
完整扩容控制器
把推理服务调用、决策算法、K8s 执行串起来:
// cmd/scaler/main.go
package main
import (
"bytes"
"context"
"encoding/json"
"fmt"
"log"
"net/http"
"os"
"os/signal"
"syscall"
"time"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/clientcmd"
"yourapp/pkg/scaler"
)
func main() {
kubeconfig := os.Getenv("KUBECONFIG")
predictorURL := os.Getenv("PREDICTOR_URL") // 如 http://qps-predictor:8080
targetDeployment := os.Getenv("TARGET_DEPLOYMENT")
targetNamespace := os.Getenv("TARGET_NAMESPACE")
if predictorURL == "" || targetDeployment == "" {
log.Fatal("PREDICTOR_URL and TARGET_DEPLOYMENT are required")
}
// K8s client
config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
if err != nil {
log.Fatalf("build config failed: %v", err)
}
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
log.Fatalf("create clientset failed: %v", err)
}
scaleConfig := scaler.DefaultConfig()
// 记录上次扩缩容时间
lastScaleTime := time.Time{} // zero value 表示从未扩缩容过
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer cancel()
ticker := time.NewTicker(60 * time.Second) // 每分钟决策一次
defer ticker.Stop()
log.Printf("scaler started, target=%s/%s, predictor=%s",
targetNamespace, targetDeployment, predictorURL)
// 首次立即执行
runScaleCycle(ctx, clientset, predictorURL, targetNamespace, targetDeployment,
&lastScaleTime, scaleConfig)
for {
select {
case <-ctx.Done():
log.Println("shutting down...")
return
case <-ticker.C:
runScaleCycle(ctx, clientset, predictorURL, targetNamespace, targetDeployment,
&lastScaleTime, scaleConfig)
}
}
}
func runScaleCycle(
ctx context.Context,
clientset *kubernetes.Clientset,
predictorURL, namespace, deployment string,
lastScaleTime *time.Time,
config scaler.ScaleConfig,
) {
// 1. 调推理服务获取预测
predictions, err := fetchPrediction(ctx, predictorURL, config.PredictSteps)
if err != nil {
log.Printf("fetch prediction failed: %v", err)
return
}
// 2. 获取当前副本数
dep, err := clientset.AppsV1().Deployments(namespace).Get(ctx, deployment, metav1.GetOptions{})
if err != nil {
log.Printf("get deployment failed: %v", err)
return
}
currentReplicas := *dep.Spec.Replicas
// 3. 决策
decision := scaler.Decide(predictions, currentReplicas, *lastScaleTime, config)
log.Printf("decision: action=%s, current=%d, desired=%d, predictedMaxQPS=%.0f, reason=%s",
decision.Action, decision.CurrentReplicas, decision.DesiredReplicas,
decision.PredictedMaxQPS, decision.Reason)
// 4. 执行
if decision.Action == "no_action" {
return
}
if decision.DesiredReplicas == currentReplicas {
return
}
// 用 Scale subresource 改副本数
scale, err := clientset.AppsV1().Deployments(namespace).GetScale(ctx, deployment, metav1.GetOptions{})
if err != nil {
log.Printf("get scale failed: %v", err)
return
}
scale.Spec.Replicas = decision.DesiredReplicas
_, err = clientset.AppsV1().Deployments(namespace).UpdateScale(ctx, deployment, scale, metav1.UpdateOptions{})
if err != nil {
log.Printf("update scale failed: %v", err)
return
}
*lastScaleTime = time.Now()
log.Printf("scaled %s/%s from %d to %d",
namespace, deployment, currentReplicas, decision.DesiredReplicas)
}
// httpClient 带超时的 HTTP 客户端,避免推理服务无响应时永久阻塞
var httpClient = &http.Client{Timeout: 30 * time.Second}
// fetchPrediction 调推理服务获取预测
// 推理服务的 FastAPI 端点从 JSON body 接收 {"steps": N}
func fetchPrediction(ctx context.Context, predictorURL string, steps int) ([]scaler.PredictedPoint, error) {
url := fmt.Sprintf("%s/predict", predictorURL)
// 构造 JSON body:FastAPI 端用 Pydantic 模型从 body 接收参数
reqBody, err := json.Marshal(map[string]int{"steps": steps})
if err != nil {
return nil, fmt.Errorf("marshal request body failed: %w", err)
}
req, err := http.NewRequestWithContext(ctx, "POST", url, bytes.NewReader(reqBody))
if err != nil {
return nil, err
}
req.Header.Set("Content-Type", "application/json")
resp, err := httpClient.Do(req)
if err != nil {
return nil, fmt.Errorf("call predictor failed: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
return nil, fmt.Errorf("predictor returned %d", resp.StatusCode)
}
var result struct {
Predictions []scaler.PredictedPoint `json:"predictions"`
}
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
return nil, fmt.Errorf("decode response failed: %w", err)
}
return result.Predictions, nil
}
踩坑提示:
QPSPerPod这个值必须通过压测得出,不能拍脑袋。我一开始设了 500,结果 Pod 在 400 QPS 时 P99 就飙到 2 秒了——Goroutine 泄漏导致的。压测后改成 200,安全多了。ScaleUpThreshold设 0.7 而不是 1.0 是关键。等 100% 满载再扩就晚了——Pod 启动需要时间,这期间流量继续涨就会过载。0.7 留 30% 余量刚好覆盖 Pod 启动延迟。CooldownSeconds不能太短。我设过 30 秒,结果 Pod 频繁扩缩容——预测 5 分钟后 QPS 高就扩,下一分钟预测低了又缩。设成 120 秒后稳定多了。- 缩容每次最多减 30%,别一次缩到底。流量预测有可能误判,缩太快会导致反复横跳。
与 HPA 的联动
预测性扩容和 HPA 不是替代关系,而是互补。我让预测器负责"大方向",HPA 负责"精细微调":
# api-gateway-hpa.yaml
# HPA 配置:预测器负责提前扩容,HPA 兜底处理突发流量
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: api-gateway-hpa
namespace: prod
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: api-gateway
minReplicas: 3 # 跟预测器的 MinReplicas 保持一致
maxReplicas: 50 # 跟预测器的 MaxReplicas 保持一致
metrics:
# 基于 CPU 利用率做精细调整
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 60 # CPU > 60% 就扩容
# 基于内存做保护
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: 80
behavior:
# 扩容行为:快扩
scaleUp:
stabilizationWindowSeconds: 0 # 扩容不等待,立刻执行
policies:
# 每次最多扩容 100%(翻倍)
- type: Percent
value: 100
periodSeconds: 15
# 或者每次最多加 4 个
- type: Pods
value: 4
periodSeconds: 15
selectPolicy: Max # 取两个策略中扩容更快的那个
# 缩容行为:慢缩,避免抖动
scaleDown:
stabilizationWindowSeconds: 300 # 缩容前等 5 分钟确认流量真的降了
policies:
# 每次最多缩 10%
- type: Percent
value: 10
periodSeconds: 60
selectPolicy: Min # 取更保守的策略
协同逻辑:
- 预测器每分钟调推理服务,如果预测未来 15 分钟 QPS 会超过 70% 容量,提前扩容。它管"提前布防"。
- HPA 持续监控实时 CPU,如果预测器没预测到突发流量(比如某个活动突然上了热搜),HPA 兜底扩容。它管"实时兜底"。
- 两者都会改 Deployment 的 replicas,但不会冲突——因为 HPA 的扩容速度比预测器快(15 秒 vs 60 秒),所以突发流量由 HPA 处理,周期性流量由预测器处理。
踩坑提示:
stabilizationWindowSeconds: 300这个缩容等待窗口非常重要。不设的话 HPA 一看到 CPU 降了就立刻缩,流量一抖就又得扩——Pod 反复创建销毁,不仅浪费资源还容易出问题。- 预测器和 HPA 同时改 replicas 不会冲突,但要注意:如果预测器刚扩到 20 副本,HPA 看到 CPU 低想缩回 10,会把预测器的扩容效果抵消掉。解决方法是让 HPA 的 minReplicas 跟预测器的输出对齐,或者干脆把 HPA 的缩容
stabilizationWindowSeconds设长一点。
MLOps 流程
模型不能训一次就不管了——业务在变,流量模式也在变。需要一套 MLOps 流水线保证模型持续有效:
┌─────────────┐ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ 数据版本化 │───▶│ 实验追踪 │───▶│ 模型注册 │───▶│ CI/CD 部署 │
│ DVC │ │ MLflow │ │ MLflow │ │ ArgoCD │
└─────────────┘ └─────────────┘ └─────────────┘ └─────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────────┐
│ 监控与漂移检测 │
│ │
│ 预测误差监控 ──▶ MAPE > 20%? ──▶ 触发重训 ──▶ 回到数据版本化 │
│ (Prometheus) (Alertmanager) (Airflow) │
└─────────────────────────────────────────────────────────────────────┘
具体实践:
- 数据版本化:用 DVC 管理从 Prometheus 拉取的训练数据,每次拉取生成一个版本。出问题时可以回溯到特定版本的数据。
- 实验追踪:MLflow 记录每次训练的超参数、特征列表、模型指标(MAE/MAPE)。对比不同实验找最优配置。
- 模型注册:最佳模型注册到 MLflow Model Registry,打上
staging/production标签。 - CI/CD 集成:代码变更或定时触发重新训练。训练通过后自动构建新镜像,推到 Registry,ArgoCD 检测到新 tag 自动部署。
- 监控与漂移检测:推理服务持续记录预测值和实际值的偏差。Prometheus 采集
prediction_mae指标,当 MAPE 连续 1 小时超过 20%,Alertmanager 触发重训。
漂移检测的 PromQL:
# 预测误差率(MAPE)
# prediction_error = abs(predicted_qps - actual_qps) / actual_qps
avg_over_time(
prediction_abs_error[1h]
) / avg_over_time(
actual_qps[1h]
) * 100 > 20
踩坑提示:
- 漂移检测的阈值别设太严。我一开始设 MAPE > 10% 就告警重训,结果模型每周重训 3 次,GPU 费用爆炸。后来改成 MAPE > 20% 持续 1 小时才重训,频率降到每月 1-2 次。
- 重训的时候别用最新一天的数据。刚过去的 24 小时可能包含异常流量(比如某次故障导致 QPS 暴跌),用这种数据训练会"学坏"。我的做法是训练数据截止到 T-1(昨天),当天数据只用来验证。
- 模型部署用金丝雀而不是直接替换。新模型先接 10% 流量,对比预测误差,确认没问题再全量切换。
总结
基于流量预测的自动扩容是 AIOps 在资源优化领域的典型应用。核心思路很简单:用历史数据训练模型预测未来流量,提前扩容。但工程落地链路很长——数据采集、特征工程、模型训练、推理服务、扩容决策、K8s 执行,每一环都有坑。
我这套系统上线后效果是实打实的:大促期间 P99 延迟从之前的 3-8 秒降到了 500ms 以内,提前扩容的成功率约 85%(剩下的 15% 靠 HPA 兜底)。资源利用率从之前的平均 30% 提升到 55%——因为不再需要为"万一流量来了"预留大量冗余副本。
但别过度迷信预测模型。预测器再准也是基于历史的统计推断,遇到黑天鹅事件(突然热搜、DDoS)照样抓瞎。HPA 兜底 + 人工告警永远不能少,预测器只是让你在大部分时间里过得舒服一点。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 原生 HPA 是基于"当前指标"做扩缩容的,为什么在 30 秒内从 500 QPS 飙到 5000 QPS 这种大促式脉冲流量下会雪崩?预测式扩容从哪个环节切断了雪崩链?
- 本文把完整的预测扩容链路拆成了哪 6 个环节?请说出每一环的职责,并指出哪几环"坑最多"。
- Prophet + XGBoost 集成的思路是什么?最终预测值是怎么算出来的?为什么残差要交给 XGBoost 而不是让 Prophet 一把梭?
- 扩容决策器
Decide里为什么要"跳过 Pod 启动延迟内的预测点"再去取最大值?ScaleUpThreshold设成 0.7 而不是 1.0,工程上是为了防什么? - 预测器和 HPA 都在改同一个 Deployment 的
replicas,为什么不会打架?两者分别负责哪类流量?
动手练习(建议真做一遍):
- 跑通
fetch_metrics.py从 Prometheus 拉 30 天 QPS,再做特征工程生成qps_features.csv,观察qps_lag_24h在特征重要性里是不是排第一;故意删掉节假日标记,看看双十一那天的预测偏差有多大。 - 本地起 FastAPI 推理服务,调
/predict拿未来 15 分钟预测 QPS;然后把预测 QPS 人为调到阈值之上,跑一遍Decide,确认它算出"提前扩容"的目标副本数,并对照DefaultConfig看余量和冷却窗口怎么生效。 - 把 HPA 的
stabilizationWindowSeconds缩到 0 跑一天,观察 Pod 是否被流量抖动带着反复横跳;再叠加预测器(提前布防),对比 P99 延迟和资源利用率的变化。
本章小结
- 预测式扩容的本质是"看预报带伞":用历史数据训模型,在峰值到来前把副本扩好,把 HPA 的"事后补救"变成"提前布防"。
- 链路有 6 环节(采集 → 特征 → 训练 → 推理 → 决策 → 执行),每一环都有工程坑,其中特征工程的周期/节假日信号和
QPSPerPod压测值最影响效果。 - 集成思路是"Prophet 抓周期基线 + XGBoost 学残差",最终预测 = 基线 + 残差;决策器靠"跳过启动延迟取最大预测 + 阈值留余量 + 冷却窗口 + 缩容限速"这四条护栏稳住。
- 预测器与 HPA 是互补而非替代:周期性流量交给预测器,突发流量交给 HPA 兜底,两者改同一
replicas不冲突。 - 模型要持续有效必须上 MLOps 闭环(数据版本化 → 实验追踪 → 模型注册 → 部署 → MAPE 漂移检测触发重训),但黑天鹅事件仍要靠 HPA + 人工告警兜底。
预测扩容解决了"流量来了才手忙脚乱"的问题,但集群里还有另一类麻烦——Pod 异常、Node 抖动这些"已经出问题"的信号,下一篇我们就用 client-go 把这些异常自动捞出来并联动 LLM 做根因分析。