Spark MLlib 实战:基于协同过滤算法构建商品推荐模型并部署上线
·
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. 上线优化策略
- 冷启动处理:
- 新用户:采用热门商品推荐
- 新商品:基于内容相似度推荐
- 增量训练:
model.setColdStartStrategy("drop") // 忽略未知用户/商品 val updatedModel = als.fit(newData.union(oldData)) - 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,启用动态分配
关键注意事项:生产环境需处理模型漂移问题,建议每周全量重训+每日增量更新,监控预测分布变化。
更多推荐
所有评论(0)