训练流量预测模型实现自动扩容

2025-11-09T14:11:02+08:00 | 24分钟阅读 | 更新于 2025-11-09T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 讲清"预测式扩容"相比原生 HPA 解决了什么问题——为什么"事后补救"会在大促式脉冲流量下雪崩。
  2. 把一条完整的预测扩容链路拆成 6 个环节(采集 → 特征 → 训练 → 推理 → 决策 → 执行),说清每环职责与典型坑。
  3. 解释 Prophet + XGBoost 集成的思路:Prophet 抓周期/节假日基线,XGBoost 学残差,最终预测 = 基线 + 残差。
  4. 看懂扩容决策器的核心公式与工程护栏(预热时间、冷却窗口、缩容幅度、阈值留余量),并知道这些参数为什么是这么设的。
  5. 设计"预测器 + HPA"的协同分工,以及用 MAPE 漂移检测触发重训的 MLOps 闭环。

前置知识

  • Kubernetes 基础:Deployment、HPA、client-go 基本认知。
  • Python 基础(Pandas、scikit-learn 概念),能读模型训练脚本。
  • Go 基础,能读决策器 / 控制器代码。
  • 上一篇笔记的 client-go 执行器(本文扩容执行器与其共用代码)。

本章你会动手做的事

  1. 用文末的 fetch_metrics.py 从 Prometheus 拉 30 天 QPS,跑通特征工程生成 qps_features.csv
  2. 本地起 FastAPI 推理服务,调 /predict 拿到未来 15 分钟预测 QPS。
  3. 跑一遍扩容决策器 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 做精细调整   │
│                        └──────────────┘                              │
└─────────────────────────────────────────────────────────────────────┘

主要组件的职责:

  1. 指标采集模块:从 Prometheus 拉取历史 QPS/CPU/内存数据,写成 CSV 或直接存入时序数据库。
  2. 特征工程模块:处理时间序列,生成训练样本。这是决定模型效果的关键步骤。
  3. 模型训练模块:Prophet 做基线预测,XGBoost 做残差修正,两个模型集成。
  4. 推理服务:FastAPI 部署,每分钟返回未来 30 分钟的预测 QPS。
  5. 扩容决策器:Go 程序,每分钟调推理服务,根据预测结果算目标副本数。
  6. 执行器:调 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)                          │
└─────────────────────────────────────────────────────────────────────┘

具体实践:

  1. 数据版本化:用 DVC 管理从 Prometheus 拉取的训练数据,每次拉取生成一个版本。出问题时可以回溯到特定版本的数据。
  2. 实验追踪:MLflow 记录每次训练的超参数、特征列表、模型指标(MAE/MAPE)。对比不同实验找最优配置。
  3. 模型注册:最佳模型注册到 MLflow Model Registry,打上 staging / production 标签。
  4. CI/CD 集成:代码变更或定时触发重新训练。训练通过后自动构建新镜像,推到 Registry,ArgoCD 检测到新 tag 自动部署。
  5. 监控与漂移检测:推理服务持续记录预测值和实际值的偏差。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 兜底 + 人工告警永远不能少,预测器只是让你在大部分时间里过得舒服一点。

自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 原生 HPA 是基于"当前指标"做扩缩容的,为什么在 30 秒内从 500 QPS 飙到 5000 QPS 这种大促式脉冲流量下会雪崩?预测式扩容从哪个环节切断了雪崩链?
  2. 本文把完整的预测扩容链路拆成了哪 6 个环节?请说出每一环的职责,并指出哪几环"坑最多"。
  3. Prophet + XGBoost 集成的思路是什么?最终预测值是怎么算出来的?为什么残差要交给 XGBoost 而不是让 Prophet 一把梭?
  4. 扩容决策器 Decide 里为什么要"跳过 Pod 启动延迟内的预测点"再去取最大值?ScaleUpThreshold 设成 0.7 而不是 1.0,工程上是为了防什么?
  5. 预测器和 HPA 都在改同一个 Deployment 的 replicas,为什么不会打架?两者分别负责哪类流量?

动手练习(建议真做一遍)

  1. 跑通 fetch_metrics.py 从 Prometheus 拉 30 天 QPS,再做特征工程生成 qps_features.csv,观察 qps_lag_24h 在特征重要性里是不是排第一;故意删掉节假日标记,看看双十一那天的预测偏差有多大。
  2. 本地起 FastAPI 推理服务,调 /predict 拿未来 15 分钟预测 QPS;然后把预测 QPS 人为调到阈值之上,跑一遍 Decide,确认它算出"提前扩容"的目标副本数,并对照 DefaultConfig 看余量和冷却窗口怎么生效。
  3. 把 HPA 的 stabilizationWindowSeconds 缩到 0 跑一天,观察 Pod 是否被流量抖动带着反复横跳;再叠加预测器(提前布防),对比 P99 延迟和资源利用率的变化。

本章小结

  • 预测式扩容的本质是"看预报带伞":用历史数据训模型,在峰值到来前把副本扩好,把 HPA 的"事后补救"变成"提前布防"。
  • 链路有 6 环节(采集 → 特征 → 训练 → 推理 → 决策 → 执行),每一环都有工程坑,其中特征工程的周期/节假日信号和 QPSPerPod 压测值最影响效果。
  • 集成思路是"Prophet 抓周期基线 + XGBoost 学残差",最终预测 = 基线 + 残差;决策器靠"跳过启动延迟取最大预测 + 阈值留余量 + 冷却窗口 + 缩容限速"这四条护栏稳住。
  • 预测器与 HPA 是互补而非替代:周期性流量交给预测器,突发流量交给 HPA 兜底,两者改同一 replicas 不冲突。
  • 模型要持续有效必须上 MLOps 闭环(数据版本化 → 实验追踪 → 模型注册 → 部署 → MAPE 漂移检测触发重训),但黑天鹅事件仍要靠 HPA + 人工告警兜底。

预测扩容解决了"流量来了才手忙脚乱"的问题,但集群里还有另一类麻烦——Pod 异常、Node 抖动这些"已经出问题"的信号,下一篇我们就用 client-go 把这些异常自动捞出来并联动 LLM 做根因分析。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!