Spark MLlib 实战:基于协同过滤算法构建商品推荐模型并部署上线

1. 问题定义与数据准备

目标:构建商品推荐系统,根据用户历史行为预测偏好
数据集:用户-商品交互数据(如评分、购买记录)

# 示例数据结构
+------+-------+-------+
|userId|itemId |rating |
+------+-------+-------+
| 1    | A     | 5.0   |
| 1    | B     | 3.5   |
| 2    | A     | 4.0   |
+------+-------+-------+

2. 协同过滤算法原理

使用交替最小二乘法 (ALS) 进行矩阵分解: $$ \underset{U,V}{\text{min}} \sum_{(i,j)\in \Omega} (r_{ij} - u_i^T v_j)^2 + \lambda (|u_i|^2 + |v_j|^2) $$ 其中:

  • $U$:用户隐因子矩阵
  • $V$:商品隐因子矩阵
  • $\lambda$:正则化系数
3. 模型构建 (Scala/Python)
import org.apache.spark.ml.recommendation.ALS

// 1. 加载数据
val ratings = spark.read.parquet("user_ratings.parquet")

// 2. 构建ALS模型
val als = new ALS()
  .setRank(10)          // 隐因子维度
  .setMaxIter(15)       // 迭代次数
  .setRegParam(0.1)     // 正则化参数 $\lambda$
  .setUserCol("userId")
  .setItemCol("itemId")
  .setRatingCol("rating")

// 3. 训练模型
val model = als.fit(ratings)

4. 模型评估

使用均方根误差 (RMSE) 评估: $$ \text{RMSE} = \sqrt{\frac{1}{N} \sum_{i=1}^{N} (y_i - \hat{y}_i)^2} $$

val predictions = model.transform(testData)
val evaluator = new RegressionEvaluator()
  .setMetricName("rmse")
  .setLabelCol("rating")
  .setPredictionCol("prediction")
val rmse = evaluator.evaluate(predictions)
println(s"模型RMSE: $rmse")

5. 生成推荐结果
// 为每个用户推荐Top-K商品
val userRecs = model.recommendForAllUsers(10) 

// 结果示例
+------+-----------------------------+
|userId| recommendations            |
+------+-----------------------------+
| 1    | [(B,4.8), (C,4.5), ...]    |
| 2    | [(D,5.0), (A,4.7), ...]    |
+------+-----------------------------+

6. 模型部署方案

方案一:批处理模式 (每日更新)

spark-submit --class RecommenderBatch \
  --master yarn \
  recommender.jar \
  --input new_ratings/ \
  --model_path /models/als

方案二:API实时服务 (Spark Serving)

# Flask API示例
from flask import Flask, request
from pyspark.sql import SparkSession

app = Flask(__name__)
spark = SparkSession.builder.getOrCreate()
model = ALSModel.load("models/als")

@app.route('/recommend', methods=['POST'])
def recommend():
    user_id = request.json['user_id']
    user_df = spark.createDataFrame([(user_id,)], ["userId"])
    recs = model.recommendForUserSubset(user_df, 10)
    return recs.toJSON().first()

7. 上线优化策略
  1. 冷启动处理
    • 新用户:采用热门商品推荐
    • 新商品:基于内容相似度推荐
  2. 增量训练
    model.setColdStartStrategy("drop")  // 忽略未知用户/商品
    val updatedModel = als.fit(newData.union(oldData))
    

  3. AB测试
    • 分组对比ALS与ItemCF的效果
    • 监控点击率(CTR)、转化率(CVR)
8. 架构设计
graph LR
A[用户行为日志] --> B[Spark ETL]
B --> C[模型训练]
C --> D[模型存储]
D --> E[API服务]
E --> F[推荐结果]
G[监控系统] --> C & E

9. 性能调优建议
  • 数据分区repartition(1000) 避免数据倾斜
  • 参数优化
    val paramGrid = new ParamGridBuilder()
      .addGrid(als.rank, Array(5, 10, 20))
      .addGrid(als.regParam, Array(0.01, 0.1, 0.5))
      .build()
    

  • 资源分配:Executor内存 ≥ 8GB,启用动态分配

关键注意事项:生产环境需处理模型漂移问题,建议每周全量重训+每日增量更新,监控预测分布变化。

更多推荐