推荐系统:协同过滤与矩阵分解的原理与 PySpark 实现

推荐系统是信息过滤技术,用于预测用户对项目的偏好(如电影、商品)。协同过滤(Collaborative Filtering)基于用户行为相似性,矩阵分解(Matrix Factorization)则通过降维建模隐含特征。下面我将逐步解释原理,并提供 PySpark 实现代码。实现基于真实算法(如 ALS),确保可靠。

1. 协同过滤的原理

协同过滤的核心是“物以类聚,人以群分”。它分为两类:

  • 基于用户的协同过滤:如果用户 A 和 B 有相似行为,则 A 喜欢的项目可能被 B 喜欢。计算用户相似度使用余弦相似度: $$ \text{sim}(u, v) = \frac{\sum_{i \in I} r_{ui} \cdot r_{vi}}{\sqrt{\sum_{i \in I} r_{ui}^2} \cdot \sqrt{\sum_{i \in I} r_{vi}^2}} $$ 其中,$u$ 和 $v$ 是用户,$r_{ui}$ 是用户 $u$ 对项目 $i$ 的评分,$I$ 是共同评分的项目集。
  • 基于项目的协同过滤:类似项目有相似评分模式。预测用户 $u$ 对项目 $i$ 的评分: $$ \hat{r}{ui} = \frac{\sum{j \in N(i)} \text{sim}(i, j) \cdot r_{uj}}{\sum_{j \in N(i)} |\text{sim}(i, j)|} $$ 其中,$N(i)$ 是与 $i$ 最相似的项目集合,$\text{sim}(i, j)$ 是项目相似度。

协同过滤的优点是简单直观,但缺点是无法处理新用户或新项目(冷启动问题)。

2. 矩阵分解的原理

矩阵分解通过降维学习用户和项目的隐含特征。假设评分矩阵 $R$(维度 $m \times n$,$m$ 用户数,$n$ 项目数)可分解为: $$ R \approx U V^T $$ 其中,$U$ 是用户特征矩阵(维度 $m \times k$),$V$ 是项目特征矩阵(维度 $n \times k$),$k$ 是隐含因子数(通常 $k \ll \min(m, n)$)。目标是最小化重构误差: $$ \min_{U,V} \sum_{(u,i) \in \Omega} (r_{ui} - \mathbf{u}_u^T \mathbf{v}_i)^2 + \lambda (|U|_F^2 + |V|_F^2) $$ 其中,$\Omega$ 是已知评分集合,$\lambda$ 是正则化系数,$|\cdot|_F$ 是 Frobenius 范数。常用算法是交替最小二乘法(ALS),它交替优化 $U$ 和 $V$。

矩阵分解能处理高维稀疏数据,并捕捉隐含特征(如用户偏好主题),但需要较多计算资源。

3. PySpark 实现

PySpark 是 Apache Spark 的 Python API,适合大规模数据处理。使用 MLlib 库的 ALS 算法实现矩阵分解。以下是完整代码示例,包括数据加载、模型训练和预测。

from pyspark.sql import SparkSession
from pyspark.ml.recommendation import ALS
from pyspark.sql.functions import col

# 初始化 Spark 会话
spark = SparkSession.builder \
    .appName("RecommendationSystem") \
    .getOrCreate()

# 示例数据:用户ID、项目ID、评分(0-5)
data = [(0, 0, 4.0), (0, 1, 2.0), (1, 0, 5.0), (1, 1, 1.0), (2, 0, 3.0), (2, 1, 4.0)]
columns = ["user_id", "item_id", "rating"]
df = spark.createDataFrame(data, columns)

# 划分训练集和测试集
train, test = df.randomSplit([0.8, 0.2])

# 创建 ALS 模型(基于矩阵分解)
als = ALS(
    maxIter=10,  # 迭代次数
    regParam=0.01,  # 正则化系数 λ
    userCol="user_id",
    itemCol="item_id",
    ratingCol="rating",
    coldStartStrategy="drop"  # 处理冷启动:忽略未知用户/项目
)

# 训练模型
model = als.fit(train)

# 预测测试集
predictions = model.transform(test)
predictions.show()

# 评估模型(使用 RMSE)
from pyspark.ml.evaluation import RegressionEvaluator
evaluator = RegressionEvaluator(metricName="rmse", labelCol="rating", predictionCol="prediction")
rmse = evaluator.evaluate(predictions)
print(f"Root Mean Squared Error (RMSE): {rmse}")

# 为用户推荐项目(例如,用户0推荐前2个项目)
user_recs = model.recommendForAllUsers(2)
user_recs.show()

# 停止 Spark 会话
spark.stop()

代码解释:

  • 数据准备:创建模拟评分数据(用户ID、项目ID、评分),实际应用中可替换为真实数据集(如 MovieLens)。
  • 模型训练ALS 类实现矩阵分解,参数包括 maxIter(迭代次数)、regParam(正则化系数 $\lambda$)。coldStartStrategy="drop" 处理冷启动。
  • 预测与评估:使用 transform 预测评分,RegressionEvaluator 计算 RMSE(均方根误差)评估模型精度。
  • 推荐功能recommendForAllUsers 为每个用户生成 top-N 推荐项目。
  • 注意事项:确保 Spark 环境配置正确;实际数据量大时,需调整分区和资源。
4. 总结
  • 协同过滤:基于相似性,适合小规模数据,但冷启动问题明显。
  • 矩阵分解:通过 ALS 等算法高效处理稀疏数据,PySpark 实现可扩展至大数据。
  • 实际应用:在电商或流媒体中,结合两者(如混合推荐)可提升精度。PySpark 代码可直接运行,建议使用分布式集群处理真实数据。

通过此实现,您可快速构建推荐系统。如有具体数据或问题,欢迎提供细节进一步优化!

更多推荐