电商 AI 基础设施的架构演进:从规则引擎到深度模型的渐进式替换路径
一、AI 基础设施替换的工程困境
电商平台经过多年发展,积累了数以千计的规则:风控规则(IP 频率、设备指纹)、搜索排序规则(销量权重、好评率、价格区间)、推荐策略(协同过滤、关联规则)。每一条规则背后是半结构化配置和定时更新的离线数据。当业务要求引入深度学习模型提升效果时,工程团队面临一个困境:全量替换风险极高,但逐模块替换则面临规则和模型并行运行的接口一致性问题。
规则引擎的优势是确定性:输入相同,输出永远一致,排查问题只需逐条审查规则。深度模型的优势是泛化能力:捕捉高维特征的非线性关系,处理规则无法穷举的长尾 case。但模型输出的不可解释性和"静默退化"(数据漂移导致精度下降)让运维团队难以信任。
渐进式替换的核心思想是"影子模式"(Shadow Mode):新模型与旧规则并行运行,模型输出仅用于评估和对照,不参与线上决策。当模型在 A/B 测试中持续优于规则 2 周以上且无异常,逐步切换流量。流水线抽象接口保证规则和模型对外暴露相同的输入/输出契约——业务方无感知。
二、渐进式替换的架构设计
渐进式替换的四个阶段:
- 第一阶段(规则为主):模型离线训练和评估,不产生线上影响。验证离线指标(AUC、NDCG)是否优于规则。
- 第二阶段(影子模式):5% 流量走模型,结果与规则对比。关键指标:模型输出与规则输出的一致性分布、模型独有错误(规则正确但模型错误)的比例。
- 第三阶段(模型为主):95% 流量走模型,5% 保留规则作为基线和降级通道。持续监控线上指标(CTR、CVR、GMV)。
- 第四阶段(全量切换):规则引擎降为降级备份——仅当模型服务异常或指标异常时自动回切。
抽象接口层的设计是替换成功的关键。所有 AI 服务对外暴露统一接口:predict(Context) -> Decision,内部实现可以是规则、机器学习模型或深度学习模型。接口层的路由逻辑根据实验配置动态选择后端实现,业务代码无需改动。
三、统一抽象接口与灰度路由的 Rust 实现
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::RwLock;
use rand::Rng;
/// 统一的推理上下文
/// 设计原因:所有决策模块共享同一输入结构
/// 新增特征只需扩展 struct 而无需修改接口
#[derive(Debug, Clone)]
pub struct PredictContext {
pub user_id: String,
pub item_id: String,
pub scene: String,
/// 扩展特征——允许各模块按需注入
pub features: HashMap<String, serde_json::Value>,
}
/// 统一的推理结果
#[derive(Debug, Clone)]
pub struct PredictDecision {
pub score: f64,
pub action: DecisionAction,
pub metadata: HashMap<String, String>,
}
#[derive(Debug, Clone)]
pub enum DecisionAction {
Recommend,
Filter,
Rank { position: usize },
Flag { risk_level: String },
}
/// AI 推理的统一特征——所有后端实现此接口
/// 设计原因:规则引擎和深度学习模型对外暴露相同契约
/// 业务代码依赖此 trait 而无需知道具体实现
#[async_trait::async_trait]
pub trait Predictor: Send + Sync {
/// 返回预测器名称——用于日志和监控
fn name(&self) -> &str;
/// 核心推理方法
async fn predict(&self, ctx: &PredictContext) -> Result<PredictDecision, PredictError>;
/// 健康检查——用于熔断和降级判断
fn health_check(&self) -> HealthStatus;
}
#[derive(Debug)]
pub struct PredictError {
pub message: String,
pub retryable: bool,
}
#[derive(Debug)]
pub enum HealthStatus {
Healthy,
Degraded { reason: String },
Unhealthy { reason: String },
}
/// 规则引擎实现
/// 设计原因:保留为基线对比和降级备份
struct RuleEngine {
rules: Vec<Box<dyn Rule>>,
name: String,
}
#[async_trait::async_trait]
impl Predictor for RuleEngine {
fn name(&self) -> &str { &self.name }
async fn predict(&self, ctx: &PredictContext) -> Result<PredictDecision, PredictError> {
// 规则链式执行:每条规则对 score 进行增量调整
let mut decision = PredictDecision {
score: 0.0,
action: DecisionAction::Recommend,
metadata: HashMap::new(),
};
for rule in &self.rules {
rule.apply(ctx, &mut decision)?;
}
Ok(decision)
}
fn health_check(&self) -> HealthStatus {
HealthStatus::Healthy // 规则引擎无状态,始终健康
}
}
trait Rule: Send + Sync {
fn apply(&self, ctx: &PredictContext, decision: &mut PredictDecision) -> Result<(), PredictError>;
}
/// 深度学习模型实现
/// 设计原因:通过 FFI 调用 PyTorch/TensorRT 推理
struct DeepModel {
name: String,
model_path: String,
}
#[async_trait::async_trait]
impl Predictor for DeepModel {
fn name(&self) -> &str { &self.name }
async fn predict(&self, ctx: &PredictContext) -> Result<PredictDecision, PredictError> {
// FFI 调用模型的推理接口
let features = Self::extract_features(ctx);
let scores = self.infer_batch(&[features]).await?;
Ok(PredictDecision {
score: scores[0],
action: DecisionAction::Recommend,
metadata: HashMap::new(),
})
}
fn health_check(&self) -> HealthStatus {
HealthStatus::Healthy
}
}
impl DeepModel {
fn extract_features(ctx: &PredictContext) -> Vec<f32> {
// 特征工程——从上下文提取数值特征
Vec::new()
}
async fn infer_batch(&self, features: &[Vec<f32>]) -> Result<Vec<f64>, PredictError> {
// FFI 调用
Ok(vec![0.0])
}
}
/// 灰度路由器——渐进式替换的核心
/// 设计原因:根据实验配置动态选择后端
/// 支持流量百分比、用户白名单、A/B 分桶
struct GrayscaleRouter {
/// 各后端的流量权重
/// 例如: [("rules", 0.95), ("model", 0.05)]
/// 权重和必须为 1.0——初始化时校验
backends: Arc<RwLock<Vec<(Arc<dyn Predictor>, f64)>>>,
/// 降级后端——当主后端 Unhealthy 时使用
fallback: Arc<dyn Predictor>,
}
impl GrayscaleRouter {
/// 路由到一个或多个后端
/// 设计原因:影子模式下,同时调用规则和模型
/// 主流量走目标后端,影子调用的结果仅记录不返回
async fn predict(
&self,
ctx: &PredictContext,
shadow_mode: bool,
) -> Result<PredictDecision, PredictError> {
let backends = self.backends.read().await;
// 基于用户 ID 哈希的确定性分桶
// 同一用户始终路由到同一后端——保证体验一致性
let hash = Self::hash_user(&ctx.user_id);
let mut cumulative = 0.0;
let mut chosen = None;
for (predictor, weight) in backends.iter() {
cumulative += weight;
if hash <= cumulative || cumulative >= 1.0 {
chosen = Some(predictor.clone());
break;
}
}
let predictor = chosen.unwrap_or_else(|| self.fallback.clone());
// 健康检查——不健康则降级
match predictor.health_check() {
HealthStatus::Unhealthy { ref reason } => {
tracing::warn!(
predictor = predictor.name(),
reason = reason.as_str(),
"predictor unhealthy, falling back"
);
return self.fallback.predict(ctx).await;
}
_ => {}
}
let result = predictor.predict(ctx).await;
// 影子模式:额外调用对比后端,仅记录差异
if shadow_mode {
self.shadow_compare(ctx, predictor.name()).await;
}
result
}
async fn shadow_compare(&self, ctx: &PredictContext, primary_name: &str) {
let backends = self.backends.read().await;
for (predictor, _) in backends.iter() {
if predictor.name() == primary_name {
continue;
}
// 影子调用——忽略错误,仅记录指标
if let Ok(shadow_result) = predictor.predict(ctx).await {
tracing::info!(
primary = primary_name,
shadow = predictor.name(),
"shadow comparison logged"
);
}
}
}
fn hash_user(user_id: &str) -> f64 {
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
let mut hasher = DefaultHasher::new();
user_id.hash(&mut hasher);
(hasher.finish() as f64) / (u64::MAX as f64)
}
/// 更新流量分配——动态调整灰度比例
/// 设计原因:线上发现问题时即时调整权重——无需重启
async fn update_weights(&self, new_weights: Vec<(Arc<dyn Predictor>, f64)>) {
let total: f64 = new_weights.iter().map(|(_, w)| w).sum();
assert!((total – 1.0).abs() < 0.001, "weights must sum to 1.0");
*self.backends.write().await = new_weights;
}
}
四、渐进式替换的风险控制与适用边界
适用场景:存量系统复杂且不可停机——规则和模型必须并行运行数月。业务对效果指标敏感——A/B 测试数据驱动切换决策。模型效果需要线上验证——离线指标(AUC)不能完全代表线上指标(CTR)。需保留快速回滚能力——切换不当可即时恢复为规则引擎。
不适用场景:新业务从零开始——直接使用模型,无需兼容规则。规则逻辑简单(< 20 条规则)——重写成本低于抽象接口层开发。模型输出已经明显优于规则——全量切换风险低。团队规模小(< 5 人)——维护双轨系统的额外成本不可承受。
Trade-offs:影子模式的双倍计算成本——每次请求调用两次推理,QPS 翻倍。但这是风险可控的代价——对比系统崩溃导致的业务损失。抽象接口层增加调用链长度——多一层间接调用(~2ns),对 P99 延迟影响可忽略。回滚能力的时间窗口由数据新鲜度决定——规则配置更新需保证与模型输出的时效性一致。
